
Introduction
Asynchronous communication allows services to send messages without waiting for a response. This is the backbone of Event-Driven Architecture and the key to building loosely coupled, resilient systems.
1. Why do we need Async Communication?
1.1 The problem of Synchronous
Sync chain: Order → Payment → Inventory → Notification
│
Vấn đề: │
├── Temporal coupling: Tất cả phải online cùng lúc
├── Latency: Total = sum(latency mỗi service)
├── Cascading failure: 1 service down → cả chain fail
└── Tight coupling: Order phải biết Payment, Inventory, ...
1.2 Async Solution
Async: Order ──publish event──▶ Message Broker
│
┌───────────────┼───────────────┐
▼ ▼ ▼
Payment Inventory Notification
(subscribe) (subscribe) (subscribe)
Ưu điểm:
├── Temporal decoupling: Services không cần online cùng lúc
├── Performance: Order trả response ngay, không chờ downstream
├── Loose coupling: Order không biết ai subscribe
└── Resilience: Message broker buffer khi consumer down
2. Message Queue — Point-to-Point
2.1 Concepts
One producer sends a message, only one consumer receives and processes:
Producer ──msg──▶ Queue ──msg──▶ Consumer
(FIFO)
Nếu có nhiều consumers → Load balancing (round-robin)
Producer ──▶ Queue ──▶ Consumer 1
──▶ Consumer 2
──▶ Consumer 3
Mỗi message chỉ được xử lý bởi MỘT consumer
2.2 RabbitMQ
┌──────────┐ ┌──────────────────────────────────┐ ┌───────────┐
│ Producer │───▶│ RabbitMQ │───▶│ Consumer │
└──────────┘ │ │ └───────────┘
│ ┌──────────┐ ┌──────────────┐ │
│ │ Exchange │───▶│ Queue │ │
│ │ (routing)│ │ (buffer msgs) │ │
│ └──────────┘ └──────────────┘ │
└──────────────────────────────────┘
Exchange Types:
├── Direct → Route by exact routing_key
├── Topic → Route by pattern (order.* , *.created)
├── Fanout → Broadcast to all bound queues
└── Headers → Route by message headers
2.3 Use Cases for Message Queue
- Task Queue: Background jobs (send email, generate report)
- Work Distribution: Divide work among many workers
- Rate Limiting: Buffer requests when downstream is slow
- Delayed Processing: Dead Letter Exchange + TTL
3. Event Streaming — Pub/Sub
3.1 Concepts
Producer publish event, multiple consumer groups subscribe and handle independently:
Consumer Group A
┌────▶ Payment Service
│
Producer ──▶ Topic ──────┤ Consumer Group B
(Order (order. ├────▶ Inventory Service
Service) created) │
│ Consumer Group C
└────▶ Notification Service
Mỗi consumer group nhận TẤT CẢ messages
Trong một group, messages được phân chia cho các instances
3.2 Apache Kafka Architecture
┌────────────────────────────────────────────────────────┐
│ Kafka Cluster │
│ │
│ Topic: order.created │
│ ┌────────────┐ ┌────────────┐ ┌────────────┐ │
│ │ Partition 0│ │ Partition 1│ │ Partition 2│ │
│ │ msg1, msg4 │ │ msg2, msg5 │ │ msg3, msg6 │ │
│ │ msg7, ... │ │ msg8, ... │ │ msg9, ... │ │
│ └────────────┘ └────────────┘ └────────────┘ │
│ │
│ Broker 1 Broker 2 Broker 3 │
│ (Leader P0) (Leader P1) (Leader P2) │
│ (Replica P1) (Replica P2) (Replica P0) │
│ │
│ ZooKeeper / KRaft (metadata management) │
└────────────────────────────────────────────────────────┘
Important features:
- Log-based: Messages are append-only, not deleted after consumption
- Retention: Retain messages by time (default 7 days) or size
- Offset: Each consumer group tracks its own offset
- Replay: Consumer can replay messages from any offset
- Ordering: Guarantees order in partition (does not guarantee cross-partition)
3.3 Topic Design
Naming convention: <domain>.<entity>.<event>
order.order.created
order.order.confirmed
order.order.cancelled
payment.payment.completed
payment.payment.failed
inventory.stock.reserved
inventory.stock.released
Partition key:
├── order_id → Đảm bảo events cùng order đi vào cùng partition → đúng thứ tự
├── customer_id → Events cùng customer ordered
└── random → Distribute đều, không đảm bảo ordering
4. Message Queue vs Event Streaming
| Criteria | Message Queue (RabbitMQ) | Event Streaming (Kafka) |
|---|---|---|
| Model | Point-to-Point (or Pub/Sub) | Pub/Sub (log-based) |
| Message lifetime | Deleted after consumed | Retained (configurable) |
| Replay | No | Yes |
| Ordering | Queue-level FIFO | Partition-level |
| Throughput | ~50K msg/s | ~1M+ msg/s |
| Consumer groups | Limited | Native support |
| Use case | Task queue, work distribution | Event streaming, audit log, analytics |
| Complexity | Low | Medium-High |
| Protocol | AMQP | Custom (Kafka protocol) |
When to choose which one?
RabbitMQ:
├── Background jobs (send email, generate PDF)
├── Request-reply pattern
├── Complex routing rules
├── Low throughput (< 50K msg/s)
└── Team muốn đơn giản, ít learning curve
Kafka:
├── Event sourcing / Event-driven architecture
├── Stream processing (real-time analytics)
├── Audit log (cần replay)
├── High throughput (> 100K msg/s)
├── Multiple consumer groups cho cùng event
└── Data pipeline (connect to data warehouse)
5. Event Schema Design
5.1 CloudEvents Specification
Standardize event format:
{
"specversion": "1.0",
"id": "evt-001-abc-def",
"source": "/services/order-service",
"type": "com.myorg.order.created",
"datacontenttype": "application/json",
"time": "2026-03-31T10:00:00Z",
"data": {
"order_id": "O-001",
"customer_id": "C-042",
"items": [
{"product_id": "P-100", "quantity": 2, "price": 250000}
],
"total": 500000,
"currency": "VND"
}
}
5.2 Schema Evolution
When the schema changes, backward compatibility is needed:
Schema Registry (Confluent / Apicurio):
├── v1: {order_id, customer_id, total}
├── v2: {order_id, customer_id, total, currency} ← thêm optional field
└── v3: {order_id, customer_id, total, currency, discount}
Quy tắc:
✅ Thêm optional field → backward compatible
✅ Thêm default value cho field mới
❌ Xoá required field → breaking change
❌ Đổi type của field → breaking change
6. Idempotency — Handles duplicate messages
6.1 Problem
Message broker guarantees at-least-once delivery → Consumer can receive message many times:
Producer ──msg──▶ Broker ──msg──▶ Consumer
│ │
│ (process OK, nhưng ACK bị mất)
│ │
└──msg──▶ Consumer (nhận lại!)
6.2 Solution: Idempotent Consumer
-- Lưu event_id đã xử lý
CREATE TABLE processed_events (
event_id VARCHAR(255) PRIMARY KEY,
processed_at TIMESTAMP DEFAULT NOW()
);
-- Trong consumer:
BEGIN;
-- Kiểm tra đã xử lý chưa
INSERT INTO processed_events (event_id) VALUES ('evt-001')
ON CONFLICT (event_id) DO NOTHING;
-- Nếu insert thành công (chưa xử lý) → xử lý business logic
IF FOUND THEN
UPDATE inventory SET stock = stock - 1 WHERE product_id = 'P-100';
END IF;
COMMIT;
7. Summary
| Pattern | When to use |
|---|---|
| Sync (REST/gRPC) | Need immediate response, query, simple request-reply |
| Message Queue | Background jobs, task distribution, rate buffering |
| Event Streaming | Event-driven, audit log, multiple consumers, replay |
| Idempotent Consumer | Always implement for every async consumer |
| Schema Registry | When the event schema needs to evolve without breaking consumers |
Next article: API Gateway Pattern — Single entry point for microservices system, centrally handling routing, authentication, and rate limiting.