การ scale WebSocket ยากกว่าการ scale HTTP API ทั่วไป เพราะ WebSocket แต่ละ connection ผูกอยู่กับ server ตัวที่รับมันไว้ HTTP request หนึ่งตัวจะไปลง replica ไหนก็ได้แล้วจบที่นั่น แต่ WebSocket เปิดค้างอยู่กับ process เดียวเป็นนาทีหรือเป็นชั่วโมง และมีแค่ process นั้นที่เขียนข้อมูลลง connection ได้ พอเริ่มรันมากกว่าหนึ่ง replica ข้อความที่เข้ามาที่ server ตัวหนึ่งจึงมักต้องไปถึง user ที่ต่ออยู่กับอีกตัว
วิธีแก้มาตรฐานคือวาง pub/sub ไว้ระหว่าง server ทุกตัว ทุก server publish ข้อความที่รับมาเข้า broker, subscribe broker ตัวเดียวกัน แล้วส่งต่อเฉพาะให้ connection ที่ตัวเองถืออยู่ บทความนี้อธิบายการออกแบบนั้น เทียบ Redis Pub/Sub, RabbitMQ และ Kafka สำหรับงานนี้ และเก็บรายละเอียดที่เจอจริงใน production คือ sticky session, reconnect storm และ presence
ทำไมการ scale WebSocket ถึงพังเมื่อมีมากกว่าหนึ่ง server
สมมติมี chat service รันอยู่ 3 pod user A ต่ออยู่กับ pod 1 ส่วน user B ต่ออยู่กับ pod 3 พอ A ส่งข้อความหา B ข้อความจะเข้า pod 1 แต่ pod 1 ไม่มี socket ของ B เพราะ map ของ connection ใน memory ของมันรู้จักแค่ user ที่ต่อกับ pod 1 เท่านั้น
ตอนมี server ตัวเดียวปัญหานี้ไม่มีทางโผล่ มันจึงมักเจอครั้งแรกหลัง scale out: ข้อความถึงบางคนแต่ไม่ถึงบางคน ขึ้นกับว่าแต่ละคนไปลง pod ไหน และการเปิด sticky session ก็ไม่ช่วยอะไร เพราะ A กับ B ต่อถูกที่อยู่แล้วทั้งคู่ ปัญหาคือ pod แต่ละตัวไม่มีทางคุยกันเองเลย
Fan-out pattern: ใช้ broker เป็นตัวกลางระหว่าง pod
วาง message broker ไว้หลัง pod ทั้งหมด แล้วเปลี่ยนความหมายของคำว่า "ส่ง":
- A ส่งข้อความผ่าน WebSocket เข้า pod 1
- Pod 1 ไม่พยายามหา B เอง แต่ publish ข้อความเข้า broker
- Broker กระจายข้อความให้ทุก pod ที่ subscribe อยู่
- แต่ละ pod เช็คว่าผู้รับต่ออยู่กับตัวเองไหม pod 3 เจอ B จึงเขียนข้อความลง socket ของ B ส่วน pod 1 กับ pod 2 ทิ้งไป
นอกจาก socket ที่ถืออยู่ pod ทุกตัวยังเป็น stateless เหมือนเดิม จะเพิ่มหรือลด replica ก็ไม่ต้องประสานอะไรกัน
ตัวอย่าง hub ใน Go ด้วย Redis Pub/Sub
Hub นี้ใช้ gorilla/websocket กับ go-redis v9 แต่ละ pod เก็บ map ของ connection ในเครื่องตัวเอง และถือ Redis subscription ไว้หนึ่งตัว
package realtime
import (
"context"
"encoding/json"
"log"
"net/http"
"sync"
"github.com/gorilla/websocket"
"github.com/redis/go-redis/v9"
)
const channel = "chat:messages"
type Message struct {
To string `json:"to"`
From string `json:"from"`
Body string `json:"body"`
}
type client struct {
userID string
send chan []byte
}
type Hub struct {
rdb *redis.Client
mu sync.RWMutex
clients map[string]map[*client]struct{} // userID -> connections on this pod
}
func NewHub(rdb *redis.Client) *Hub {
return &Hub{rdb: rdb, clients: make(map[string]map[*client]struct{})}
}
// Run holds one subscription per pod and delivers to local connections only.
func (h *Hub) Run(ctx context.Context) {
sub := h.rdb.Subscribe(ctx, channel)
defer sub.Close()
ch := sub.Channel()
for {
select {
case <-ctx.Done():
return
case msg, ok := <-ch:
if !ok {
return
}
var m Message
if err := json.Unmarshal([]byte(msg.Payload), &m); err != nil {
continue
}
h.deliverLocal(m.To, []byte(msg.Payload))
}
}
}
func (h *Hub) deliverLocal(userID string, payload []byte) {
h.mu.RLock()
defer h.mu.RUnlock()
for c := range h.clients[userID] {
select {
case c.send <- payload:
default: // client is too slow; never block the whole pod
}
}
}
func (h *Hub) add(c *client) {
h.mu.Lock()
defer h.mu.Unlock()
if h.clients[c.userID] == nil {
h.clients[c.userID] = make(map[*client]struct{})
}
h.clients[c.userID][c] = struct{}{}
}
func (h *Hub) remove(c *client) {
h.mu.Lock()
defer h.mu.Unlock()
delete(h.clients[c.userID], c)
if len(h.clients[c.userID]) == 0 {
delete(h.clients, c.userID)
}
close(c.send)
}
// The default CheckOrigin rejects cross-origin upgrades.
var upgrader = websocket.Upgrader{}
func (h *Hub) ServeWS(w http.ResponseWriter, r *http.Request) {
userID := r.Header.Get("X-User-ID") // set by your auth middleware
if userID == "" {
http.Error(w, "unauthorized", http.StatusUnauthorized)
return
}
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
return
}
defer conn.Close()
c := &client{userID: userID, send: make(chan []byte, 64)}
h.add(c)
defer h.remove(c)
// gorilla/websocket allows one concurrent writer per connection.
go func() {
for payload := range c.send {
if err := conn.WriteMessage(websocket.TextMessage, payload); err != nil {
return
}
}
}()
for {
_, data, err := conn.ReadMessage()
if err != nil {
return
}
var m Message
if err := json.Unmarshal(data, &m); err != nil {
continue
}
m.From = userID // never trust a sender field from the client
out, err := json.Marshal(m)
if err != nil {
continue
}
if err := h.rdb.Publish(r.Context(), channel, out).Err(); err != nil {
log.Printf("publish: %v", err)
}
}
}มีสองจุดที่สำคัญ อย่างแรก ห้ามให้ loop ของ subscriber ถูก block ถ้า client ช้าตัวเดียวทำให้ deliverLocal ค้างได้ ทุกคนบน pod นั้นจะไม่ได้ข้อความไปด้วย โค้ดจึงใช้ buffered channel คู่กับ default เพื่อทิ้งข้อความแทน (ใน production มักตัด connection ของ client ตัวนั้นไปเลย) อย่างที่สอง บันทึกข้อความก่อน publish เสมอ ประวัติแชทต้องอยู่ใน database ส่วน pub/sub มีหน้าที่แค่ส่งสด
ออกแบบ channel: กระจายทุกอย่าง หรือ subscribe เฉพาะที่สนใจ
ตัวอย่างข้างบนส่งทุกข้อความไปทุก pod ซึ่งพอใช้ได้ถ้ามีไม่กี่ pod แต่ pod ส่วนใหญ่จะทิ้งข้อความส่วนใหญ่ที่ได้รับ และงานที่เสียเปล่าจะโตตามจำนวน pod คูณจำนวนข้อความ ขั้นถัดไปคือ channel ตามความสนใจ: publish ไปที่ room:<id> หรือ user:<id> แล้วให้แต่ละ pod subscribe เฉพาะ channel ของห้องหรือ user ที่ตัวเองถือ connection อยู่ go-redis เรียก Subscribe และ Unsubscribe บน PubSub ตัวเดิมได้ตามจังหวะที่ connection เข้าออก แลกกับ bookkeeping เพิ่มขึ้นนิดหน่อย แต่ traffic ที่เสียเปล่าลดลงมาก
Redis Pub/Sub vs RabbitMQ vs Kafka สำหรับ realtime fan-out
ทั้งสามตัวทำงานนี้ได้ ต่างกันที่การรับประกันการส่งและต้นทุนในการดูแล
| Redis Pub/Sub | RabbitMQ | Kafka | |
|---|---|---|---|
| การส่ง | At-most-once, fire-and-forget | มี ack, ตั้งให้ durable ได้ | เก็บเป็น log, replay ได้ |
| ข้อความที่ publish ตอน pod หลุดการเชื่อมต่อ | หาย | อยู่ ถ้า queue ยังอยู่ | อยู่จนหมด retention |
| Latency | ต่ำสุด | ต่ำ | สูงกว่า เพราะจูนมาเพื่อ throughput |
| การตั้งค่า fan-out | subscriber ทุกตัวของ channel ได้ทุกข้อความ | fanout หรือ topic exchange, หนึ่ง queue ต่อ pod | หนึ่ง consumer group ต่อ pod |
| ภาระการดูแล | เบา | ปานกลาง | หนักสุด |
| เหมาะกับ | แชท, notification, live dashboard | ข้อความที่ห้ามหาย, routing ซับซ้อน | ปริมาณมหาศาล, event ชุดเดียวป้อนหลายระบบ |
Redis Pub/Sub เป็นตัวเลือกตั้งต้นที่นิยมที่สุด และ Redis adapter ของ Socket.IO ก็สร้างบนมัน ข้อแลกคือ Redis ไม่เก็บข้อความ ถ้า pod หลุดจาก Redis ไปชั่วครู่ ข้อความที่ publish ในช่วงนั้นจะหายไปสำหรับ pod นั้น สำหรับแชทสดมักรับได้ เพราะตอน reconnect client จะโหลดประวัติล่าสุดจาก database อยู่แล้ว ถ้าต้องการ buffer ให้ใช้ Redis Streams ซึ่งมี persistence และ consumer group บน server เดิม อีกเรื่องที่ควรรู้คือบน Redis Cluster คำสั่ง PUBLISH แบบเดิมจะถูกกระจายไปทุก node ส่วน Redis 7 เพิ่ม sharded pub/sub (SPUBLISH และ SSUBSCRIBE) ที่ให้แต่ละ channel อยู่บน shard ของตัวเอง อ่านต่อได้ที่ Redis caching, data structure และ persistence
RabbitMQ เหมาะเมื่อข้อความ realtime ห้ามหาย เช่นการยืนยันการชำระเงิน แต่ละ pod ประกาศ queue ของตัวเองแบบ exclusive และ auto-delete แล้ว bind เข้ากับ fanout หรือ topic exchange ถ้าใช้ topic exchange ก็ route ตามห้องหรือตาม tenant ได้ด้วย รายละเอียดอยู่ใน RabbitMQ exchange และการส่งข้อความให้เชื่อถือได้
Kafka คุ้มเมื่อ event ชุดเดียวกันต้องป้อนทั้ง analytics, audit และ search ด้วย แต่มีกับดักหนึ่งข้อสำหรับงาน fan-out consumer ที่อยู่ group เดียวกันจะแบ่ง partition กันอ่าน ถ้า WebSocket pod ทุกตัวใช้ group ID เดียวกัน แต่ละ pod จะเห็นข้อความแค่บางส่วน ถ้าต้องการ broadcast ให้แต่ละ pod มี group ของตัวเองและเริ่มอ่านจาก offset ล่าสุด เหตุผลอธิบายไว้ใน Apache Kafka partition และ consumer group
Sticky session: เมื่อไหร่ต้องใช้ เมื่อไหร่ไม่ต้อง
WebSocket แท้ ๆ ไม่ต้องใช้ sticky session เพราะหลัง upgrade แล้ว connection จะอยู่กับ pod เดิมจนกว่าจะปิด
กรณีที่ต้องใช้จริง ๆ คือ:
- Client fallback ไปใช้ HTTP long-polling ได้ Socket.IO และ SockJS อาจส่ง HTTP request หลายตัวแยกกันที่เป็นของ session เดียว ถ้า request เหล่านั้นไปลงคนละ pod การ handshake จะล้ม ทางแก้คือเปิด affinity (ด้วย cookie หรือ client IP) ที่ load balancer หรือตั้งให้ client ใช้ WebSocket transport อย่างเดียว
- เก็บ session state ที่ต้อง resume ไว้ใน memory ของ pod sticky ช่วยได้ แต่ทางที่ดีกว่ามักเป็นการย้าย state นั้นไปไว้ใน Redis
อย่าลืมเช็ค idle timeout ด้วย proxy และ cloud load balancer หลายตัวปิด connection ที่เงียบเกินประมาณหนึ่งนาที ให้เพิ่ม timeout (ใน ingress-nginx คือ proxy-read-timeout และ proxy-send-timeout) หรือส่ง ping ถี่กว่า timeout ที่สั้นที่สุดบนเส้นทาง
Rollout, การ scale out และ reconnect storm
Connection ที่อยู่ยาวทำให้ deployment และ autoscaling ทำงานต่างจาก HTTP ปกติ:
- Pod ใหม่ไม่ได้รับ connection เดิมไปด้วย scale จาก 3 เป็น 6 replica แล้ว socket เดิมก็ยังอยู่ที่เดิม มีแค่ connection ใหม่ที่กระจายไป pod ใหม่
- ทุกครั้งที่ pod ปิด socket ทั้งหมดของมันหลุด rolling update หนึ่งรอบอาจทำให้ client หลายพันตัว reconnect พร้อมกัน client ควร retry แบบ exponential backoff และสุ่ม jitter ส่วน server ควรรับ
SIGTERMแล้วส่ง close frame และทยอยปิด ไม่ใช่ตายทันที - Scale ตามจำนวน connection ไม่ใช่แค่ CPU socket ที่ว่างอยู่กิน memory และ file descriptor แต่แทบไม่กิน CPU
ติดตาม presence ข้าม pod
คำถามว่า "B ออนไลน์อยู่ไหม" มีปัญหาแบบเดียวกับการส่งข้อความ คือคำตอบกระจายอยู่หลาย pod จึงควรเก็บ presence ไว้ใน Redis ไม่ใช่ใน memory ของ pod
Sorted set หนึ่งตัวต่อ user ใช้ได้ดี แต่ละ connection ใส่ ID ของตัวเองโดยใช้ Unix timestamp ปัจจุบันเป็น score แล้วอัปเดตทุกครั้งที่ส่ง heartbeat ถ้ามี entry ไหนที่ยังใหม่อยู่ แปลว่า user ออนไลน์
# on connect and on every heartbeat (for example every 30s)
ZADD presence:user:42 1790000000 conn-7f3a
# online if any connection checked in within the last 60s
ZCOUNT presence:user:42 1789999940 +inf
# on clean disconnect
ZREM presence:user:42 conn-7f3a
# periodic cleanup of entries left by crashed pods
ZREMRANGEBYSCORE presence:user:42 -inf 1789999940การหมดอายุตาม heartbeat สำคัญเพราะ pod ที่ crash จะไม่ได้รัน handler ตอน disconnect เลย และการนับ connection แทนการเก็บ flag ตัวเดียว ยังรองรับ user ที่เปิดหลาย tab ได้ด้วย
เวลาแจ้งสถานะให้คนอื่น ให้ publish การเปลี่ยน presence เฉพาะคนที่เกี่ยวข้อง เช่นเพื่อนหรือสมาชิกในห้อง และทำ debounce ไว้ เพื่อไม่ให้มือถือที่สัญญาณไม่นิ่งสร้าง event ออนไลน์/ออฟไลน์รัว ๆ
คำถามที่พบบ่อย
Server หนึ่งตัวรับ WebSocket connection ได้กี่ connection
ขึ้นกับ workload ของคุณ จึงต้อง load test เอง ข้อจำกัดที่มักเจอก่อนคือ file descriptor, memory สำหรับ buffer และ proxy ที่อยู่ข้างหน้า ไม่ใช่ CPU
WebSocket ต้องใช้ sticky session ไหม
ถ้าเป็น WebSocket แท้ไม่ต้อง เพราะ connection อยู่กับ server เดียวอยู่แล้ว จะต้องใช้ก็ต่อเมื่อ library ของคุณ fallback ไปเป็น HTTP long-polling ได้ แบบที่ Socket.IO ทำ
Redis Pub/Sub เชื่อถือได้พอสำหรับแชทไหม
สำหรับการส่งสด ส่วนใหญ่พอ ถ้าบันทึกข้อความลง database ก่อน และ client โหลดประวัติล่าสุดตอน reconnect แต่ถ้าข้อความ realtime หายไม่ได้เลย ให้ใช้ Redis Streams, RabbitMQ หรือ Kafka
ใช้ Kafka แทน Redis ในการ scale WebSocket ได้ไหม
ได้ แต่ต้องให้ WebSocket pod แต่ละตัวมี consumer group ของตัวเอง ทุก pod จะได้เห็นทุกข้อความ ถ้าใช้แค่ fan-out อย่างเดียว Kafka มักเกินความจำเป็น
Checklist สำหรับการ scale WebSocket
- ทุก pod publish ข้อความที่รับมาเข้า broker และส่งต่อเฉพาะ connection ของตัวเอง
- บันทึกข้อความลง database ก่อน publish
- Client ที่ช้าไม่สามารถ block loop ของ subscriber ได้
- Idle timeout ของ proxy ยาวกว่ารอบ ping
- Client reconnect แบบมี backoff และ jitter และ pod ทยอยปิดตอน shutdown
- Presence อยู่ใน Redis และหมดอายุตาม heartbeat ไม่ได้อยู่ใน memory ของ pod
- เปิด sticky session เฉพาะเมื่อมี long-polling fallback
เริ่มจาก Redis Pub/Sub กับ channel เดียว วัดผล แล้วค่อยขยับไปใช้ channel ตามความสนใจหรือ broker ที่ durable เมื่อตัวเลขหรือความต้องการด้านการส่งบังคับให้ต้องเปลี่ยน ถ้ากำลังวางแผนฟีเจอร์ realtime และอยากได้ความเห็นอีกมุมเรื่องสถาปัตยกรรม ทีม Vectorkub ช่วยออกแบบและรีวิวระบบแบบนี้ได้
