Node.js Production Engineering 19 — Thiết kế Realtime Systems
Chọn polling, SSE hay WebSocket từ product contract; triển khai delivery, reconnect, authentication, presence, backpressure, scale-out và graceful draining với Node.js.
Một dashboard hiển thị “đơn hàng đã thanh toán”. Người dùng mất mạng ba giây đúng lúc event được phát. Kết nối tự reconnect, giao diện có chấm xanh trở lại — nhưng trạng thái vẫn là “đang xử lý” vì server không lưu event đã bỏ lỡ.
Kết nối lại thành công không đồng nghĩa khôi phục trạng thái thành công. Realtime production là bài toán delivery, ordering, recovery và capacity trước khi là bài toán chọn thư viện WebSocket.
Sau bài này, bạn sẽ có thể:
- viết realtime contract từ nhu cầu sản phẩm thay vì mặc định WebSocket;
- chọn polling, SSE, WebSocket hoặc Socket.IO theo hướng dữ liệu và recovery;
- triển khai SSE/WebSocket có cleanup, size limit, schema validation và backpressure;
- giải thích ordering và delivery guarantee thực sự của Socket.IO;
- scale nhiều instance mà hiểu lúc nào cần sticky session và adapter;
- thiết kế reconnect, presence, token expiry, deploy draining và observability.
1. Bắt đầu bằng realtime contract
Trước công nghệ, ghi rõ sáu câu hỏi:
- Hướng dữ liệu: chỉ server → client hay hai chiều?
- Freshness: 100 ms, 5 giây hay 1 phút vẫn chấp nhận?
- Delivery: có được mất event không? Duplicate có gây hại không?
- Ordering: theo một document/user hay toàn hệ thống?
- Recovery: reconnect cần replay event hay chỉ refetch snapshot mới nhất?
- Scale: bao nhiêu kết nối đồng thời, message/giây và bytes/giây?
Ví dụ:
Stock ticker freshness < 1s, server→client, bỏ intermediate update được
Order status không được kẹt trạng thái, reconnect có thể refetch snapshot
Collaborative doc hai chiều, ordering theo document, cần conflict/recovery model
Presence chấp nhận eventual, TTL, disconnect không luôn đáng tin
Hai feature đều “realtime” nhưng reliability contract hoàn toàn khác. Đừng bắt chat, price tick và payment status dùng cùng một delivery policy chỉ vì cùng đi qua socket.
2. Chọn transport nhỏ nhất đáp ứng contract
| Cơ chế | Hướng | Điểm mạnh | Chi phí/giới hạn | Hợp với |
|---|---|---|---|---|
| Polling | client → server định kỳ | đơn giản, cache/proxy quen thuộc | request thừa, freshness theo interval | trạng thái đổi chậm |
| Long polling | server giữ request tới khi có dữ liệu | tương thích HTTP rộng | nhiều request lifecycle | fallback/legacy |
| SSE | server → client | HTTP text stream, browser reconnect, event id | một chiều, text, browser API hạn chế header | notification, feed, progress |
| WebSocket | hai chiều | một kết nối duplex, overhead message thấp | tự lo protocol, auth lifecycle, recovery | chat, collaboration, game |
| Socket.IO | hai chiều + abstraction | room, ack, reconnect, fallback, adapters | không phải raw WebSocket protocol, thêm state/overhead | ứng dụng cần feature framework |
Nếu client chỉ cần biết trạng thái mới nhất và độ trễ 10 giây chấp nhận được, polling có conditional request có thể là lựa chọn tốt hơn một connection sống nhiều giờ. Nếu chỉ server push, SSE thường ít moving part hơn WebSocket.
3. SSE đúng: stream, resume và proxy behavior
Native EventSource tự reconnect và gửi Last-Event-ID khi server cung cấp id. Muốn resume, server vẫn phải giữ event log hoặc map cursor tới dữ liệu bền.
import { once } from 'node:events';
import type { Request, Response } from 'express';
app.get('/events', async (req: Request, res: Response) => {
const actor = await authenticateRequest(req);
const after = req.get('last-event-id') ?? undefined;
const abort = new AbortController();
res.writeHead(200, {
'Content-Type': 'text/event-stream; charset=utf-8',
'Cache-Control': 'no-cache, no-transform',
Connection: 'keep-alive',
'X-Accel-Buffering': 'no', // Nginx: đừng buffer stream
});
res.flushHeaders();
const heartbeat = setInterval(() => {
res.write(': heartbeat\n\n'); // comment SSE, giữ proxy connection sống
}, 15_000);
// Cleanup theo response/socket đóng, không theo request body đã đọc xong.
res.once('close', () => {
clearInterval(heartbeat);
abort.abort();
});
try {
for await (const event of eventStore.subscribe({
actorId: actor.id,
after,
signal: abort.signal,
})) {
const frame = [
`id: ${event.id}`,
`event: ${event.type}`,
`data: ${JSON.stringify(event.data)}`,
'',
'',
].join('\n');
if (!res.write(frame)) await once(res, 'drain');
}
} catch (error) {
if (!abort.signal.aborted) logger.error({ error }, 'SSE stream failed');
}
});
res.write() trả false nghĩa là buffer đã vượt high-water mark; chờ drain là backpressure. Heartbeat trong code minh họa chưa phối hợp drain riêng; ở throughput cao, gom mọi write qua một connection writer để không có hai producer ghi đồng thời.
Native EventSource trong browser không cho đặt arbitrary Authorization header. Thường dùng cookie HttpOnly cùng kiểm tra Origin, hoặc signed URL sống rất ngắn — không đặt access token dài hạn trong query vì URL dễ rơi vào log/history.
4. Raw WebSocket với ws: boundary phải phòng thủ
WebSocket bắt đầu bằng HTTP upgrade rồi trở thành message channel dài hạn. Browser gửi Origin; server phải allowlist để chống cross-site WebSocket hijacking khi auth dựa trên cookie.
import { WebSocket, WebSocketServer, type RawData } from 'ws';
import { z } from 'zod';
const incomingMessage = z.discriminatedUnion('type', [
z.object({
type: z.literal('chat.send'),
roomId: z.string(),
body: z.string().max(2_000),
}),
z.object({
type: z.literal('typing.set'),
roomId: z.string(),
active: z.boolean(),
}),
]);
const wss = new WebSocketServer({
noServer: true,
maxPayload: 64 * 1024,
perMessageDeflate: false,
});
const actorBySocket = new WeakMap<WebSocket, Actor>();
httpServer.on('upgrade', async (request, socket, head) => {
try {
assertAllowedOrigin(request.headers.origin);
const actor = await authenticateUpgrade(request);
wss.handleUpgrade(request, socket, head, (ws) => {
actorBySocket.set(ws, actor);
wss.emit('connection', ws, request);
});
} catch {
socket.write('HTTP/1.1 401 Unauthorized\r\nConnection: close\r\n\r\n');
socket.destroy();
}
});
wss.on('connection', (ws: WebSocket) => {
const actor = actorBySocket.get(ws);
if (!actor) return ws.close(1011, 'missing connection context');
let chain = Promise.resolve();
let queued = 0;
async function handleIncoming(raw: RawData): Promise<void> {
let input: unknown;
try {
input = JSON.parse(raw.toString());
} catch {
return ws.close(1008, 'invalid JSON');
}
const parsed = incomingMessage.safeParse(input);
if (!parsed.success) return ws.close(1008, 'invalid message');
if (!(await roomPolicy.canPublish(actor.id, parsed.data.roomId))) {
return ws.close(1008, 'forbidden');
}
await handleMessage(actor, parsed.data);
}
ws.on('message', (raw) => {
if (queued >= 32) return ws.close(1013, 'message backlog exceeded');
queued += 1;
// EventEmitter không await async listener. Nối promise để giữ thứ tự
// theo connection, bắt rejection và giới hạn backlog trong bộ nhớ.
chain = chain
.then(() => handleIncoming(raw))
.catch((error) => {
logger.error({ error, actorId: actor.id }, 'WebSocket message failed');
ws.close(1011, 'message processing failed');
})
.finally(() => {
queued -= 1;
});
});
});
Đoạn trên rút gọn để tập trung boundary. Bản production còn cần rate limit theo actor/connection, heartbeat, structured error, trace context và shutdown. Không bao giờ tin userId, role hoặc room membership trong payload client.
perMessageDeflate có lợi với payload lớn/lặp, nhưng dùng CPU/memory và có rủi ro amplification; đo trước khi bật. maxPayload mặc định của thư viện có thể lớn hơn business cần rất nhiều, nên đặt rõ.
5. Socket.IO: protocol ứng dụng có sẵn primitive
Socket.IO thêm event, ack, room, reconnect và transport fallback. Nó không phải raw WebSocket server; client WebSocket thuần không nói chuyện trực tiếp với Socket.IO protocol.
interface ClientToServerEvents {
'chat:send': (
input: { roomId: string; clientMessageId: string; body: string },
ack: (
result: { ok: true; eventId: string } | { ok: false; code: string }
) => void
) => void;
}
interface ServerToClientEvents {
'chat:message': (event: ChatEvent) => void;
}
interface SocketData {
actor: Actor;
}
const io = new Server<
ClientToServerEvents,
ServerToClientEvents,
{},
SocketData
>(httpServer, {
cors: { origin: allowedOrigins, credentials: true },
maxHttpBufferSize: 64 * 1024,
connectionStateRecovery: {
maxDisconnectionDuration: 2 * 60 * 1000,
skipMiddlewares: false,
},
});
io.use(async (socket, next) => {
try {
socket.data.actor = await authenticateHandshake(socket.handshake);
next();
} catch {
next(new Error('unauthorized'));
}
});
io.on('connection', (socket) => {
socket.on('chat:send', async (input, ack) => {
if (!(await roomPolicy.canPublish(socket.data.actor.id, input.roomId))) {
return ack({ ok: false, code: 'FORBIDDEN' });
}
const { event, created } = await chat.appendIdempotently({
...input,
actorId: socket.data.actor.id,
});
if (created) io.to(`room:${input.roomId}`).emit('chat:message', event);
ack({ ok: true, eventId: event.id });
});
});
clientMessageId cho phép server deduplicate khi client retry vì không nhận ack. Khi store trả event cũ, server không broadcast lại; recipient vẫn deduplicate theo event.id để phòng duplicate từ các đường delivery khác. Ack chỉ chứng minh server callback đã trả, không mặc định chứng minh mọi recipient đã render message.
6. Ordering và delivery: ghi rõ phạm vi
Socket.IO đảm bảo thứ tự event đã đến trên một connection, kể cả lúc transport upgrade. Mặc định message arrival là at-most-once: connection đứt giữa lúc gửi có thể mất event, và server không tự lưu mọi event cho client đang offline.
Để server → client có khả năng replay:
1. Event có id tăng theo stream/document.
2. Server persist event trước khi emit.
3. Client chỉ lưu offset sau khi apply thành công.
4. Reconnect gửi offset cuối.
5. Server replay event sau offset hoặc yêu cầu refetch snapshot.
6. Client deduplicate theo event id.
io.on('connection', async (socket) => {
const offset = socket.handshake.auth.offset as string | undefined;
for (const event of await chat.eventsAfter(offset, { limit: 1_000 })) {
socket.emit('chat:message', event);
}
});
Connection state recovery của Socket.IO có thể phục hồi room/data và missed packet trong một cửa sổ cấu hình, nhưng recovery không phải lúc nào cũng thành công. Nó hoạt động với in-memory adapter mặc định; Redis Pub/Sub adapter chuẩn không hỗ trợ feature này. Client vẫn cần fallback full resync.
Ordering toàn cục rất đắt và hiếm khi cần. Thường chỉ cần ordering theo roomId/document/aggregate. Khi nhiều producer ghi cùng stream, dùng sequence/version do một owner cấp hoặc conflict-resolution model phù hợp.
7. Scale-out: load balancing và message forwarding là hai việc khác nhau
Khi có nhiều instance:
client A ─▶ instance 1 ─┐
├─ adapter/pub-sub ─▶ broadcast tới room ở mọi instance
client B ─▶ instance 2 ─┘
Adapter chuyển event giữa instance; nó không mặc định lưu lịch sử cho offline recovery.
import { createAdapter } from '@socket.io/redis-adapter';
import { createClient } from 'redis';
const pubClient = createClient({ url: process.env.REDIS_URL });
const subClient = pubClient.duplicate();
await Promise.all([pubClient.connect(), subClient.connect()]);
io.adapter(createAdapter(pubClient, subClient));
Sticky session: khi nào cần?
- Socket.IO còn dùng HTTP long-polling: các HTTP request của cùng session phải tới instance đã tạo session → cần sticky session hoặc cơ chế đồng bộ session tương đương.
- Chỉ dùng WebSocket/WebTransport: một connection sống ở một instance → round-robin handshake được, không cần sticky session cho các frame sau.
Tắt polling giảm compatibility ở mạng/proxy hạn chế. Đây là trade-off cần test với người dùng thật, không phải tối ưu mặc định.
Redis Pub/Sub adapter không hỗ trợ Connection State Recovery hay replay khi Redis/service ngắt. Nếu recovery là requirement, dùng Redis Streams adapter tương thích hoặc event store riêng và kiểm ma trận feature/version của adapter.
8. Presence là lease, không phải boolean tuyệt đối
Disconnect event không luôn tới: laptop ngủ, NAT timeout hoặc process chết. Presence production thường là lease có TTL:
online nếu heartbeat/lease được renew trong 30s
away nếu quá 30s
offline nếu quá 2 phút
Một user có nhiều tab/device, nên userId → boolean không đủ. Lưu nhiều connection/session lease và suy ra trạng thái user. Presence là eventual; UI nên diễn đạt “hoạt động gần đây” thay vì hứa chính xác tuyệt đối.
Không ghi một Redis write cho mọi heartbeat của hàng triệu connection mà chưa capacity plan. Có thể aggregate theo process, dùng sorted set/bucket hoặc chỉ bật presence cho room đang quan sát.
9. Backpressure và slow consumer
Nếu server tạo 1.000 message/giây nhưng client đọc 100, buffer tăng cho tới khi memory hoặc latency vỡ. Policy tùy loại dữ liệu:
- coalesce: chỉ giữ price/progress mới nhất;
- drop: bỏ typing indicator cũ;
- batch: gom analytics update trong 100 ms;
- disconnect: đóng client quá chậm để bảo vệ hệ thống;
- persist + resume: chat/order event không được mất.
Với ws, theo dõi bufferedAmount; với stream API, tôn trọng write()/drain. Đặt per-connection outbound queue hữu hạn và metric slow-consumer disconnect. “Gửi thành công vào buffer” không đồng nghĩa client đã nhận.
10. Security lifecycle của kết nối dài
Handshake auth chỉ là thời điểm đầu. Access token có thể hết hạn, user bị khóa hoặc permission room thay đổi trong khi socket vẫn sống.
Baseline:
- TLS (
wss://) và allowlistOrigin; - cookie
SameSite/CSRF model hoặc token sống ngắn trong handshake; - re-auth/disconnect khi token hết hạn hay session bị revoke;
- authorization trên mỗi loại event/resource, không chỉ lúc connect;
- schema, size, rate và concurrency limit;
- không phản chiếu error nội bộ/stack;
- audit event nhạy cảm;
- quota connection theo account/IP có cân nhắc NAT.
CORS config của Socket.IO không thay thế authorization và không bảo vệ raw WebSocket theo cách nhiều người kỳ vọng. Origin validation ở handshake mới xử lý browser cross-site case; non-browser client có thể giả Origin, nên vẫn cần authentication thực.
11. Heartbeat, reconnect và deploy draining
Heartbeat phát hiện half-open connection. Reconnect dùng exponential backoff + jitter để tránh hàng trăm nghìn client cùng kết nối lại sau deploy.
Khi deploy:
- instance chuyển
ready = falseđể load balancer không nhận connection mới; - gửi signal “server draining” hoặc đóng với code/reason đã quy ước;
- client reconnect có jitter sang instance khỏe;
- chờ connection cũ đóng tới deadline rồi terminate;
- client resume bằng offset hoặc refetch snapshot.
Nếu đóng toàn fleet cùng lúc, reconnect storm có thể hạ bản release mới. Rolling/canary và connection drain là phần của deployment plan realtime.
12. Hợp đồng với frontend
Team cần thống nhất một event envelope:
type EventEnvelope<T> = {
id: string;
type: string;
version: number;
occurredAt: string;
streamId: string;
sequence: number;
data: T;
};
Và state machine client:
connecting → live → reconnecting → recovering → live
└─ recovery failed → resyncing → live
UI không nên chỉ có isConnected. recovering/stale giúp tránh hiển thị dữ liệu cũ như dữ liệu live. Mutation gửi qua socket cần clientMessageId, ack timeout và UX cho trạng thái pending/failed.
13. Capacity và observability
Theo dõi:
- active connection theo instance/region/transport;
- connect, auth failure, reconnect và recovery success rate;
- inbound/outbound message + bytes/giây;
- event processing/ack latency;
- per-connection buffered bytes, slow-consumer drop/disconnect;
- room size và adapter publish latency;
- event replay gap, oldest unrecovered offset;
- event-loop lag, RSS và file descriptor usage.
Connection count không đủ. 50.000 connection idle có thể rẻ hơn 5.000 connection phát payload lớn 20 lần/giây. Load test phải mô phỏng message distribution, reconnect storm và slow client, không chỉ mở socket.
14. Checklist trước khi ship
- Product contract ghi freshness, delivery, ordering và recovery.
- Transport là lựa chọn nhỏ nhất đáp ứng contract.
- Message có schema/version/size limit; authorization theo resource.
- Có event id/offset hoặc chiến lược full resync khi reconnect.
- Duplicate và out-of-order được test.
- Backpressure policy phân biệt event được drop và event phải persist.
- Sticky session chỉ bật khi transport/session model cần.
- Adapter, event store và recovery responsibility được tách rõ.
- Token expiry/revocation và Origin validation có test.
- Deploy drain và reconnect storm có runbook/metric.
15. Bài thực hành: order status + support chat
Dùng hai feature để buộc hai contract khác nhau:
- Order status bằng SSE: event id,
Last-Event-ID, heartbeat, cleanup trênres.once('close', ...), Nginx buffering off. - Support chat bằng Socket.IO: room authorization,
clientMessageId, ack, event persistence và offset replay. - Chạy hai instance với Redis adapter; test polling có/không sticky session và WebSocket-only.
- Kill connection giữa lúc gửi; refresh tab; tắt Redis; deploy rolling. Ghi expected state/recovery cho từng case.
- Tạo slow client và chứng minh outbound queue bị giới hạn.
- Dashboard recovery success, oldest gap, buffered bytes và reconnect rate.
Bài đạt khi UI không nói “live” trong lúc đang recover và không mất business state sau một disconnect có chủ đích.
Tài liệu chính thức
- HTML Standard: Server-sent events
- Node.js HTTP API
- ws documentation
- Socket.IO: Delivery guarantees
- Socket.IO: Connection state recovery
- Socket.IO: Using multiple nodes
Phần tiếp theo
Realtime và microservices tạo nhiều failure path mà log rời rạc không thể giải thích. Phần cuối xây observability từ SLI/SLO: structured logs, metrics có cardinality kiểm soát, distributed traces, sampling và OpenTelemetry Collector — rồi dùng chúng trong một incident workflow hoàn chỉnh.