Chương 05 · Async Layer

Message Queue & Async Communication

Sync vs Async, Queue vs Pub/Sub, Kafka/RabbitMQ/SQS, delivery guarantees, idempotency, dead letter queue, backpressure, event-driven architecture.

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.

Client ──► A ──► B ──► C (blocking through) Tổng latency = sum của mọi hop Bất kỳ service nào fail → toàn chuỗi fail

1.2. Async communication

A push message vào queue → trả client ngay. B/C lấy message từ queue và xử lý sau.

Client ──► A ──► [QUEUE] │ ▼ (B lấy khi rảnh) B ──► [QUEUE] │ ▼ C Client không chờ B chậm/chết → message vẫn an toàn trong queue Có thể replay khi B recover

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.

Producer ──► [Queue: msg1, msg2, msg3, msg4] ──► Consumer A ──► Consumer B ──► Consumer C msg1 → A msg2 → B (mỗi msg đến đúng 1 consumer) msg3 → C msg4 → 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.

Publisher ──► [Topic: msg1] ──► Subscriber A nhận msg1 ──► Subscriber B nhận msg1 ──► Subscriber C nhận msg1 Mỗi sub độc lập, có offset riêng, có thể replay 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).

[Topic: order.created] │ ┌──────────────┴──────────────┐ ▼ ▼ ┌─ Group: email-service ─┐ ┌─ Group: analytics ─┐ │ ▶ consumer-1 (msg1) │ │ ▶ consumer-1 (msg1) │ │ ▶ consumer-2 (msg2) │ │ ▶ consumer-2 (msg2) │ └────────────────────────┘ └─────────────────────┘ (1 group: queue chia tải) (giữa groups: broadcast)

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.
TOPIC: order.events (3 partitions) Partition 0: [m1] [m4] [m7] [m10] ... ← messages append-only Partition 1: [m2] [m5] [m8] [m11] ... Partition 2: [m3] [m6] [m9] [m12] ... Consumer group "billing": Consumer-1 → Partition 0 Consumer-2 → Partition 1 Consumer-3 → Partition 2 (Mỗi partition đọc bởi 1 consumer trong group)

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?

KafkaRabbitMQAWS SQS
ModelCommit log, persistentTraditional queue + topicManaged queue
Throughput1M+ msg/s/broker50k msg/s~3000 msg/s/queue (FIFO)
Latency2-10ms< 1ms10-100ms
RetentionDays/weeks/monthsCho đến khi consumer ack14 ngày max
ReplayCó (rewind offset)Không nativeKhông
RoutingĐơn giản (partition key)Mạnh (exchange, binding, headers)Cơ bản
Operational complexityCao (Zookeeper/KRaft, JVM tuning)Trung bìnhZero (managed)
Use caseEvent streaming, log, analytics, high throughputTask queue, RPC, complex routingAWS-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

GuaranteeMô tảTrade-off
At-most-onceMessage 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-onceMessage delivered ≥ 1 lần. Có thể duplicate.Mặc định nhiều system. Consumer phải idempotent.
Exactly-onceMessage 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ông INCR).
  • "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.

Main Queue ──► Consumer ──► Process ▲ │ │ ▼ Error └─── Retry (3 times) │ ▼ Sau 3 retry [DLQ] ─► Manual investigation

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.

Order Service publish "order.created" │ ▼ ┌──────── [Event bus / Kafka] ────────┐ │ │ ▼ ▼ ▼ ▼ ▼ Email Inventory Shipping Analytics Audit svc svc svc svc svc (each subscribe và xử lý độc lập)

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)
TensePast — "OrderCreated"Imperative — "CreateOrder"
DirectionBroadcast — bất cứ ai quan tâmTargeted — gửi cho 1 service cụ thể
CouplingLooseTight (sender biết receiver)
FailureSubscriber tự handleSender expect response

EDA chuẩn dùng event, không command — đó là sức mạnh của decoupling.

10. Bài tập

  1. Setup Kafka local (Docker). Tạo topic order.events với 3 partition. Producer Node.js publish 1000 event với key = userId. Verify cùng user vào cùng partition.
  2. 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.
  3. 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.
  4. 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.
  5. 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.
  6. 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.
  7. 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:

  • User login
  • Real-time chat
  • Background work (email, transcode video, generate report), decouple service, buffer spike
  • Search
User cần response ngay → sync. Background task không cần response ngay (gửi email, generate PDF, video transcode) → async. Decouple: A publish event, B/C/D subscribe. Queue làm buffer cho burst traffic.

Kafka khác queue truyền thống ở:

  • Kafka chậm hơn
  • Kafka là distributed commit log — message persistent theo retention, có offset, replay được; không xóa khi consumer ack
  • Kafka không support multiple consumer
  • Kafka chỉ in-memory
Queue truyền thống (RabbitMQ, SQS): xóa message khi consumer ack. Kafka: append-only log, giữ message theo retention (days/weeks). Consumer track offset của mình, có thể rewind. Cùng topic phục vụ nhiều consumer group độc lập.

Kafka guarantee message order:

  • Toàn topic
  • Cross broker
  • Không có guarantee
  • Trong cùng partition; cross-partition không guarantee — chọn partition key sao cho cùng entity vào cùng partition
Order trong partition. Cùng userId → hash same partition → mọi event của user đó in-order. Nếu partition theo random → mất order. Quy tắc: key theo entity cần preserve order (userId, orderId).

"At-least-once delivery + idempotent consumer" là pattern phổ biến vì:

  • Exactly-once cross-system rất khó implement; consumer idempotent đủ giải bài toán dedup ở app layer
  • Nhanh nhất
  • Tốn ít RAM
  • Không cần code
Network không có exactly-once thuần. Kafka có "EOS" trong Kafka transaction nhưng cross-system (Kafka → DB) vẫn cần idempotency. Pattern thực tế: at-least-once + dedup qua idempotency key (DB) → "effectively-once".

Idempotency key implementation:

  • Random UUID mỗi request
  • Cookie
  • Client gửi key (header), server lưu lần xử lý đầu tiên + response, request sau với cùng key trả response cũ
  • Không cần
Stripe API pattern: client tạo idempotency key (UUID), gửi qua header. Server INSERT ON CONFLICT DO NOTHING — nếu đã tồn tại, trả response cũ. Tránh double-charge khi user click submit 2 lần hoặc retry timeout.

Dead Letter Queue dùng để:

  • Tăng throughput
  • Giữ message fail nhiều lần (poison) để debug, tránh retry vô tận block các message khác
  • Encryption
  • Cache
1 message format sai, code bug → retry mãi → block toàn queue. DLQ: sau N retry (3-5), push sang queue riêng. Alert on-call. Giữ retention dài (14+ ngày). Có replay tool sau khi fix bug.

Backpressure xảy ra khi:

  • Producer chậm
  • Network unstable
  • DB full
  • Producer phát nhanh hơn consumer xử lý → queue phình → memory/disk full nếu không kiểm soát
Strategies: bounded queue (block producer khi đầy), reject new (503), drop oldest, sample, autoscale consumer. Reactive streams cho phép consumer "pull" số message họ kham được. Monitor consumer lag để detect sớm.

Event-Driven Architecture event vs command:

  • Event = "đã xảy ra" (past tense), broadcast loose coupling; Command = "làm việc này" (imperative), targeted tight coupling
  • Event và command như nhau
  • Event chỉ Kafka, command chỉ RabbitMQ
  • Event là deprecated
Event "OrderCreated" — broadcast, ai quan tâm subscribe; producer không biết consumer. Command "CreateOrder" — targeted, sender biết receiver, expect response. EDA chuẩn dùng event để decouple thực sự.

Hoàn thành Chương 5. Tiếp theo: Chương 6 — Microservices vs Monolith →