iT邦幫忙

2026 iThome 鐵人賽

DAY 17
0
Software Development

手刻 Redis:用 Go 從零打造高效能高併發的記憶體資料庫系列 第 17

Day 17:發布與訂閱(Pub/Sub)廣播機制在 Go 中的設計與實作

  • 分享至 

  • xImage
  •  

Pub/Sub (發布/訂閱) 拿來做聊天室或事件廣播很方便,發送者不用管有誰在聽,只要往頻道一丟,訂閱的客戶端就全都會收到。

今天來實作一個thread-safe 的 Pub/Sub 管理器,重點是不要讓訂閱者清不掉,也不要讓 publish 卡住整個 server。


避開循環進口 (Circular Dependency) 的設計

在 Go 語言中,專案模組之間不允許循環進口。如果:

  • command 模組引用了 db(因為命令需要呼叫 KV 引擎)。
  • db 模組需要發布訊息,引用了 command.Client

這會導致編譯失敗。
所以我先把邊界切清楚:Redis 的 Pub/Sub 機制與資料庫 KV 引擎是分開的。 發布出去的訊息只會即時推給連線中的 client,不會寫回 DB。

因此,PubSubManager 放在 command 模組裡,只跟 Client 和頻道名稱互動。這樣可以避開 commanddb 互相引用的問題。


核心設計與實作

我在 code/command/pubsub.go 中設計了全域管理器 PubSubManager

1. 管理器結構

我用一個嵌套 map,將頻道名稱映射到一個「客戶端連線集合」:

type PubSubManager struct {
	mu       sync.RWMutex
	channels map[string]map[*Client]struct{} // channel -> set of clients
}

2. 訂閱與退訂操作

當客戶端訂閱一個頻道時,我們將其加入該頻道的 map 集合中:

func (ps *PubSubManager) Subscribe(client *Client, channel string) {
	ps.mu.Lock()
	defer ps.mu.Unlock()

	clients, exists := ps.channels[channel]
	if !exists {
		clients = make(map[*Client]struct{})
		ps.channels[channel] = clients
	}
	clients[client] = struct{}{}
	client.Subscribe(channel) // 更新客戶端內部的訂閱列表
}

3. 訊息發布與安全推送

其實剛寫 Publish 的時候,我圖省事直接在全域讀鎖裡面跑網路 Write。後來自己寫個腳本並發測試一下,發布訊息的時候整個伺服器直接卡死...才趕快改成先複製一份 Client 列表出來再發送。

發布訊息時,我讀取頻道中的所有客戶端。為了避免在網路 I/O 寫入時長時間霸佔全域讀鎖,我先複製一份客戶端列表,釋放鎖後再依次進行網路推送:

func (ps *PubSubManager) Publish(channel string, message []byte) int {
	ps.mu.RLock()
	clients, exists := ps.channels[channel]
	if !exists || len(clients) == 0 {
		ps.mu.RUnlock()
		return 0
	}

	// 複製客戶端清單,減小鎖粒度
	clientList := make([]*Client, 0, len(clients))
	for client := range clients {
		clientList = append(clientList, client)
	}
	ps.mu.RUnlock()

	// 依照 RESP 規範封裝訊息:*3\r\n$9\r\nmessage\r\n$<ch_len>\r\n<channel>\r\n$<msg_len>\r\n<message>\r\n
	pubReply := resp.NewArray([]resp.Value{
		resp.NewBulkString([]byte("message")),
		resp.NewBulkString([]byte(channel)),
		resp.NewBulkString(message),
	})
	replyBytes := pubReply.Marshal()

	count := 0
	for _, client := range clientList {
		client.mu.Lock()
		_, err := client.Conn.Write(replyBytes) // 安全寫入各自的 TCP socket
		client.mu.Unlock()
		if err == nil { count++ }
	}
	return count
}

資源清理與連線中斷防漏

如果訂閱者突然中斷連線,我們必須將它從所有頻道的訂閱列表中清除,否則會導致 map 的記憶體洩漏,並且 Master 會不斷嘗試對一個已關閉的 socket 寫入資料而報錯。

我們在 handleConnection 客戶端關閉的 defer 區塊中,加入了自動退訂的機制:

// 在 server.go 連線處理中
defer command.GlobalPubSubManager.UnsubscribeAll(clientCtx)

這樣 client 異常斷線時,Pub/Sub 裡留下的訂閱狀態也能一起被清掉。


跑起來看看

這裡先用簡化版行為測一下指令路徑。專案還沒有完整內嵌 Lua 引擎,但如果有實作 EVAL,大概可以這樣送:

# 啟動 Server
go run ./code/main.go
# EVAL "return 1" 0
printf "*3\r\n\$4\r\nEVAL\r\n\$8\r\nreturn 1\r\n\$1\r\n0\r\n" | nc localhost 6379
# 預期回覆::1 (代表回傳整數 1)

這裡主要是確認指令可以走到 handler;真正的 Lua 執行環境還得另外補。


總結

今天把 Pub/Sub 的發布訂閱流程補上,也處理了循環依賴和鎖住網路 I/O 的問題。這段實作最大的坑是無窮迴圈讀 channel 的地方,沒設好 timeout 直接讓 goroutine 暴增,後來加上 context 取消機制才把 leak 解掉。

明天來看 Pipeline,順便確認前面寫的 TCP Server 到底有沒有天然支援批次命令,大家明天見!


上一篇
Day 16:Redis 事務(Transaction)原理與 MULTI, EXEC, DISCARD 實作
系列文
手刻 Redis:用 Go 從零打造高效能高併發的記憶體資料庫17
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言