Pub/Sub (發布/訂閱) 拿來做聊天室或事件廣播很方便,發送者不用管有誰在聽,只要往頻道一丟,訂閱的客戶端就全都會收到。
今天來實作一個thread-safe 的 Pub/Sub 管理器,重點是不要讓訂閱者清不掉,也不要讓 publish 卡住整個 server。
在 Go 語言中,專案模組之間不允許循環進口。如果:
command 模組引用了 db(因為命令需要呼叫 KV 引擎)。db 模組需要發布訊息,引用了 command.Client。這會導致編譯失敗。
所以我先把邊界切清楚:Redis 的 Pub/Sub 機制與資料庫 KV 引擎是分開的。 發布出去的訊息只會即時推給連線中的 client,不會寫回 DB。
因此,PubSubManager 放在 command 模組裡,只跟 Client 和頻道名稱互動。這樣可以避開 command 和 db 互相引用的問題。
我在 code/command/pubsub.go 中設計了全域管理器 PubSubManager:
我用一個嵌套 map,將頻道名稱映射到一個「客戶端連線集合」:
type PubSubManager struct {
mu sync.RWMutex
channels map[string]map[*Client]struct{} // channel -> set of clients
}
當客戶端訂閱一個頻道時,我們將其加入該頻道的 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) // 更新客戶端內部的訂閱列表
}
其實剛寫 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 到底有沒有天然支援批次命令,大家明天見!