1. Vì sao Async Communication?
1.1. Sync (synchronous) communication
Client gọi service A → A gọi B → A chờ B response → A trả client. Nếu B chậm, client cũng chậm.
1.2. Async communication
A push message vào queue → trả client ngay. B/C lấy message từ queue và xử lý sau.
1.3. Khi nào dùng async?
- Background work: gửi email, generate PDF, thumbnail image, video transcode.
- Decouple service: A không cần biết B exist, chỉ publish event.
- Buffer spike: queue hấp thụ burst, consumer xử lý đều đặn.
- Retry/replay: persistent queue cho phép replay khi consumer fix bug.
- Fan-out: 1 event → nhiều consumer độc lập.
1.4. Khi KHÔNG dùng async?
- User cần response ngay (login, search, get profile).
- Logic transaction strict (chuyển tiền phải atomic).
- Latency-sensitive real-time (game, stock trading).
1.5. Trade-off
- (+) Resilience: 1 service chết không kéo cả hệ.
- (+) Decouple: thêm consumer mới không động producer.
- (+) Burst handling: queue làm buffer.
- (−) Eventual consistency: kết quả không có ngay.
- (−) Complexity: queue, consumer, monitoring, error handling.
- (−) Debug: stack trace cross-service khó.
2. Queue vs Pub/Sub
2.1. Queue (point-to-point)
1 message → 1 consumer xử lý. Sau khi consumer ack, message bị xóa.
Use case: job queue (gửi email, video transcode). Mỗi job xử lý 1 lần.
2.2. Pub/Sub (broadcast)
1 message → mọi subscriber nhận. Mỗi subscriber có queue riêng.
Use case: event-driven — order.created → email service + analytics + warehouse + audit log đều cần biết.
2.3. Consumer Group — best of both
Kafka concept: trong cùng topic, multiple "consumer group", mỗi group có nhiều consumer. Trong 1 group, message được phân chia (queue behavior). Giữa các group, message được broadcast (pub/sub behavior).
Hệ quả: cùng topic phục vụ cả 2 pattern. Đó là sức mạnh Kafka.
3. Kafka Deep Dive
3.1. Architecture
Kafka khác queue truyền thống: là distributed commit log. Producer append, consumer đọc theo offset.
- Topic — logical category (user.events, order.events).
- Partition — chia 1 topic thành N partition để parallelize.
- Broker — Kafka server, lưu các partition.
- Producer — publish message vào topic.
- Consumer — đọc từ topic, track offset.
- Consumer Group — consumer cùng group share work.
3.2. Partition key
Producer chọn partition: partitionId = hash(key) % numPartitions. Cùng key → cùng partition → preserve order.
// Order events theo user_id để cùng user cùng partition
producer.send({
topic: 'order.events',
messages: [{
key: String(userId),
value: JSON.stringify(orderData),
}],
});
Quy tắc: Kafka guarantee order trong partition, KHÔNG cross partition. Cần order theo user → partition theo user_id.
3.3. Offset & retention
Consumer track offset của mình. Có thể "rewind" để replay event.
# Kafka retention: giữ message bao lâu
log.retention.hours=168 # 7 ngày
log.retention.bytes=1073741824 # hoặc 1GB
Khác queue truyền thống xóa khi ack — Kafka giữ message theo retention. Producer-consumer hoàn toàn decouple.
3.4. Replication
Mỗi partition có replica trên broker khác. Leader handle read/write, followers replicate. Leader chết → follower promote.
Configure: replication.factor=3, min.insync.replicas=2 — phổ biến nhất cho prod.
3.5. Producer reliability
await producer.send({
topic: 'order.events',
messages: [...],
acks: 'all', // chờ all replicas ack (most durable)
// acks: 1 — chỉ leader ack (faster)
// acks: 0 — fire and forget (no guarantee)
});
3.6. Use cases
- Event streaming (LinkedIn-original use case).
- Activity log (Twitter, Uber click stream).
- Metrics aggregation.
- CDC (Debezium → Kafka).
- Decouple microservice.
- Real-time analytics.
4. Kafka vs RabbitMQ vs SQS — chọn cái nào?
| Kafka | RabbitMQ | AWS SQS | |
|---|---|---|---|
| Model | Commit log, persistent | Traditional queue + topic | Managed queue |
| Throughput | 1M+ msg/s/broker | 50k msg/s | ~3000 msg/s/queue (FIFO) |
| Latency | 2-10ms | < 1ms | 10-100ms |
| Retention | Days/weeks/months | Cho đến khi consumer ack | 14 ngày max |
| Replay | Có (rewind offset) | Không native | Không |
| Routing | Đơn giản (partition key) | Mạnh (exchange, binding, headers) | Cơ bản |
| Operational complexity | Cao (Zookeeper/KRaft, JVM tuning) | Trung bình | Zero (managed) |
| Use case | Event streaming, log, analytics, high throughput | Task queue, RPC, complex routing | AWS-native, simple queue, low ops |
4.1. Quy tắc chọn
- Kafka: event sourcing, > 100k msg/s, multiple consumer groups, replay history.
- RabbitMQ: complex routing (priority queue, fanout, topic exchange), task queue điển hình, low latency.
- SQS: trên AWS, đơn giản, không muốn manage broker.
- Redis Streams / Pub-Sub: cực đơn giản, đã có Redis sẵn, không cần broker mới.
- NATS: cloud-native, lightweight, < 1ms latency.
- Postgres LISTEN/NOTIFY hoặc bảng outbox: nhỏ, đã có DB, không cần broker mới.
4.2. RabbitMQ exchange types
# Direct: route theo exact routing key
queue.bind(exchange='direct', binding_key='order.created')
# Topic: pattern routing key (* = 1 word, # = nhiều)
queue.bind(exchange='topic', binding_key='order.*') # match order.created, order.cancelled
queue.bind(exchange='topic', binding_key='*.payment.*')
# Fanout: broadcast tất cả binding queue
queue.bind(exchange='fanout')
# Headers: route theo message header
5. Delivery Guarantees
5.1. Ba mức đảm bảo
| Guarantee | Mô tả | Trade-off |
|---|---|---|
| At-most-once | Message delivered 0 hoặc 1 lần. Có thể mất. | Nhanh nhất, không retry. Cho metrics, log không quan trọng từng message. |
| At-least-once | Message delivered ≥ 1 lần. Có thể duplicate. | Mặc định nhiều system. Consumer phải idempotent. |
| Exactly-once | Message delivered đúng 1 lần. | Khó nhất. Yêu cầu transactional producer + consumer. |
5.2. Vì sao "Exactly-once" khó?
Producer send → broker chết trước ack → producer retry → 2 lần message. Network không có "exactly-once" guarantee thuần — chỉ có "at-least-once + dedup".
Kafka có "exactly-once semantics" (EOS) bằng:
- Idempotent producer — producer ID + sequence number, broker dedup.
- Transactional producer — atomic across multiple partition.
- Read-process-write trong cùng Kafka transaction.
Nhưng: EOS chỉ trong Kafka. Cross-system (Kafka → DB) vẫn cần idempotency.
5.3. Pattern thực tế: At-least-once + Idempotent consumer
Đa số production dùng pattern này. Consumer tự đảm bảo dedup.
async function processOrderEvent(event: OrderEvent) {
// Idempotency key
const exists = await db.query(
'SELECT 1 FROM processed_events WHERE event_id = $1',
[event.id]
);
if (exists.rowCount > 0) {
logger.info('Already processed', event.id);
return;
}
// Process trong transaction
await db.transaction(async (tx) => {
await processBusiness(tx, event);
await tx.query(
'INSERT INTO processed_events (event_id, processed_at) VALUES ($1, NOW())',
[event.id]
);
});
}
6. Idempotency — bắt buộc cho async
Idempotent: gọi nhiều lần kết quả như gọi 1 lần. f(f(x)) = f(x).
6.1. Vì sao quan trọng?
- Consumer crash giữa "process xong" và "ack" → message redeliver → process 2 lần.
- Network timeout → producer retry → 2 message.
- User click "submit" 3 lần → 3 request.
6.2. Cách thiết kế idempotent
6.2.a. Idempotency key
// Request có header Idempotency-Key
app.post('/charge', async (req, res) => {
const idemKey = req.header('Idempotency-Key');
if (!idemKey) return res.status(400).send('idempotency key required');
const result = await db.query(
`INSERT INTO charges (idem_key, amount, status)
VALUES ($1, $2, 'pending')
ON CONFLICT (idem_key) DO NOTHING
RETURNING *`,
[idemKey, req.body.amount]
);
if (result.rowCount === 0) {
// Đã tồn tại — trả response cũ
const existing = await db.query('SELECT * FROM charges WHERE idem_key=$1', [idemKey]);
return res.json(existing.rows[0]);
}
// Mới — process
await processCharge(result.rows[0]);
res.json(result.rows[0]);
});
6.2.b. Conditional update
-- ❌ Không idempotent
UPDATE accounts SET balance = balance + 100 WHERE id = 5;
-- ✓ Với version
UPDATE accounts SET balance = balance + 100, version = version + 1
WHERE id = 5 AND version = $expectedVersion;
-- Retry không tăng nhiều lần vì version đổi
6.2.c. Upsert
INSERT INTO orders (id, ...) VALUES ($1, ...)
ON CONFLICT (id) DO NOTHING;
-- Hoặc DO UPDATE SET ... WHERE excluded.updated_at > orders.updated_at;
6.2.d. Natural idempotency
Một số operation tự idempotent:
SET,DELETE(chứ khôngINCR).- "Set status = 'paid'" (chứ không "increment paid count").
- HTTP PUT (chứ không POST).
6.3. Idempotency key storage
Lưu key trong DB hoặc Redis với TTL (1-7 ngày). Tránh table phình mãi.
7. Dead Letter Queue (DLQ)
Message fail nhiều lần (poison message) — không retry vô tận, đẩy sang queue riêng để debug.
7.1. Setup pattern
async function processWithRetry(msg: Message) {
const retryCount = msg.headers['retry-count'] || 0;
try {
await process(msg);
await ack(msg);
} catch (err) {
if (retryCount >= 3) {
// Đẩy sang DLQ
await publishToDLQ({
...msg,
headers: { ...msg.headers, error: err.message, original_topic: msg.topic },
});
await ack(msg);
} else {
// Retry với backoff
await publishWithDelay({
...msg,
headers: { ...msg.headers, 'retry-count': retryCount + 1 },
}, exponentialBackoff(retryCount));
await ack(msg);
}
}
}
7.2. DLQ best practices
- Alert: DLQ > 0 → ping on-call.
- Investigate: error type, frequency, root cause.
- Replay: sau khi fix, có tool replay DLQ về main queue.
- Retention dài: giữ DLQ ≥ 14 ngày để debug.
7.3. Exponential backoff
function exponentialBackoff(attempt: number, base = 1000) {
const exp = Math.pow(2, attempt) * base; // 1s, 2s, 4s, 8s, 16s
const jitter = Math.random() * exp * 0.3; // ±30% jitter
return Math.min(exp + jitter, 60_000); // cap 60s
}
Jitter quan trọng — tránh "retry storm" khi nhiều consumer cùng retry sau cùng thời điểm.
8. Backpressure — controlling flow
Khi producer phát nhanh hơn consumer xử lý → queue phình → memory/disk full → crash.
8.1. Strategies
- Bounded queue — set max size. Producer block khi queue đầy.
- Reject new messages — return 503 cho client. Họ retry sau.
- Drop oldest — discard message cũ (cho metrics, log không critical).
- Sample — chỉ ghi 10% message khi rate cao.
- Scale consumer — autoscale theo queue depth.
8.2. Reactive Streams
Backpressure protocol — consumer "pull" số lượng message họ sẵn sàng xử lý:
// Pseudocode reactive stream
consumer.request(10); // tôi xử lý được 10 message
producer.send(10);
// consumer xử lý xong
consumer.request(20); // request thêm
Implemented trong: RxJS, Project Reactor, Akka Streams. Native trong gRPC streaming.
8.3. Monitor lag
- Queue depth — số message chưa xử lý.
- Consumer lag — offset cuối producer - offset cuối consumer.
- Processing rate — msg/s consumer xử lý.
- Producer rate — msg/s producer phát.
Khi lag tăng dần → consumer không kịp. Scale up hoặc fix bottleneck.
9. Event-Driven Architecture (EDA)
9.1. Khái niệm
Service không gọi trực tiếp nhau — communicate qua event. Service A publish event "order.created" — bất cứ service nào quan tâm đều subscribe.
9.2. Lợi ích
- Loose coupling — Order không biết Email/Inventory exist.
- Easy add consumer — service mới subscribe topic, không động producer.
- Scalable — mỗi service scale độc lập.
- Resilient — 1 service chết, queue hold message.
- Auditable — event stream là audit log tự nhiên.
9.3. Trade-off
- Eventual consistency — kết quả không có ngay.
- Debug khó — flow distributed, stack trace cross-service.
- Schema evolution — đổi event schema = đụng tất cả consumer.
- Operational complexity — Kafka, monitoring, alerting.
9.4. Schema registry
Confluent Schema Registry quản lý schema (Avro/JSON Schema/Protobuf) để producer-consumer agree. Compatibility mode (BACKWARD, FORWARD, FULL) cho phép evolve schema an toàn.
9.5. Event vs Command
| Event (đã xảy ra) | Command (yêu cầu làm) | |
|---|---|---|
| Tense | Past — "OrderCreated" | Imperative — "CreateOrder" |
| Direction | Broadcast — bất cứ ai quan tâm | Targeted — gửi cho 1 service cụ thể |
| Coupling | Loose | Tight (sender biết receiver) |
| Failure | Subscriber tự handle | Sender expect response |
EDA chuẩn dùng event, không command — đó là sức mạnh của decoupling.
10. Bài tập
- Setup Kafka local (Docker). Tạo topic
order.eventsvới 3 partition. Producer Node.js publish 1000 event với key = userId. Verify cùng user vào cùng partition. - Implement consumer group "billing" với 3 consumer. Mỗi consumer xử lý 1 partition. Kill 1 consumer giữa chừng — quan sát rebalancing.
- Implement at-least-once delivery + idempotent consumer cho "process payment":
- Idempotency key.
- Save processed event ID vào DB.
- Consumer crash giữa chừng → message redelivered → không process duplicate.
- So sánh tools: chọn message broker phù hợp cho 4 use case:
- (a) Audit log — high throughput, retain 30 ngày.
- (b) Task queue gửi email — low volume, complex routing.
- (c) AWS-native simple queue.
- (d) Cross-microservice event-driven, replay được.
- Implement DLQ pattern: 3 retry với exponential backoff + jitter. Sau 3 fail → push DLQ. Tạo replay tool đọc DLQ → push lại main queue.
- Backpressure: producer phát 10k msg/s, consumer xử lý 5k msg/s. Sau 1h chuyện gì xảy ra? 3 cách giải.
- Thiết kế event schema cho domain e-commerce:
order.created,order.paid,order.shipped,order.cancelled. Xác định: producer, các consumer, payload field.
11. Quiz
Quiz cuối Chương 5
Async qua message queue phù hợp NHẤT cho:
Kafka khác queue truyền thống ở:
Kafka guarantee message order:
"At-least-once delivery + idempotent consumer" là pattern phổ biến vì:
Idempotency key implementation:
Dead Letter Queue dùng để:
Backpressure xảy ra khi:
Event-Driven Architecture event vs command:
Hoàn thành Chương 5. Tiếp theo: Chương 6 — Microservices vs Monolith →