Chuyển đến nội dung chính

Lesson 13: Message Queues & Task Queues

Why is asynchronous processing needed? Message Queue vs Task Queue. RabbitMQ architecture & exchange types. Apache Kafka architecture. Queue patterns: Work Queue, Pub/Sub, Dead Letter Queue. Guaranteed message delivery & idempotency.

🏗️ Architecture — Lesson 13 Lesson 13: Message Queues & Task Queues

System Architecture: From Zero to Hero

Part 4: Asynchronous Processing & Communication

xdev.asia

Introduction

Synchronous processing: User clicks "Order" → waits 30 seconds (create order, send email, calculate shipping fee, update inventory...). Too slow! Message Queues allows async processing, reduces latency and increases reliability.


1. Synchronous vs Asynchronous

1.1 The problem of Synchronous

User: POST /api/orders
  │
  ├── [200ms] Validate order
  ├── [100ms] Charge payment
  ├── [300ms] Update inventory
  ├── [500ms] Send confirmation email
  ├── [200ms] Notify warehouse
  ├── [150ms] Update analytics
  │
  └── Response: 1.45 giây! 😱

Nếu Email Service down → Toàn bộ request fail!

1.2 Asynchronous with Message Queue

User: POST /api/orders
  │
  ├── [200ms] Validate order
  ├── [100ms] Charge payment
  ├── [5ms]   Publish "OrderCreated" event
  │
  └── Response: 305ms! 🚀

Background workers xử lý async:
  Queue → [Worker 1] Update inventory
  Queue → [Worker 2] Send email
  Queue → [Worker 3] Notify warehouse
  Queue → [Worker 4] Update analytics

2. Message Queue vs Task Queue

Message Queue:                    Task Queue:
┌─────────────────┐              ┌─────────────────┐
│ Producer gửi    │              │ Producer gửi    │
│ MESSAGE (data)  │              │ TASK (function   │
│                 │              │  + arguments)    │
│ Consumer tự     │              │ Worker thực      │
│ quyết định      │              │ thi task         │
│ xử lý thế nào  │              │                  │
│                 │              │ Retry, schedule  │
│ RabbitMQ, Kafka │              │ Celery, Sidekiq  │
│ SQS, NATS       │              │ Bull, Temporal   │
└─────────────────┘              └─────────────────┘

3. RabbitMQ

3.1 AMQP Architecture

Producer → Exchange → Binding → Queue → Consumer

Exchanges types:
┌──────────────────────────────────────────────────────┐
│ Direct Exchange:                                     │
│   Message với routing_key="email"                    │
│   → Queue "email_queue" (binding key="email")        │
│                                                      │
│ Fanout Exchange:                                     │
│   Message broadcast đến TẤT CẢ queues bound         │
│   → Queue A, Queue B, Queue C (all get copy)         │
│                                                      │
│ Topic Exchange:                                      │
│   Message routing_key="order.created.vn"             │
│   → Queue 1 bound "order.created.*" ✅               │
│   → Queue 2 bound "order.#" ✅                       │
│   → Queue 3 bound "payment.*" ❌                     │
│                                                      │
│ Headers Exchange:                                    │
│   Route dựa trên message headers thay vì routing_key │
└──────────────────────────────────────────────────────┘

3.2 Work Queue Pattern

               ┌── Consumer 1 (xử lý message 1, 3, 5)
               │
Producer ──► Queue ──── Consumer 2 (xử lý message 2, 4, 6)
               │
               └── Consumer 3 (xử lý message 7, 8, 9)

Mỗi message chỉ được 1 consumer xử lý
Load balancing: Round-robin hoặc prefetch

3.3 Pub/Sub Pattern

                    ┌── Queue A ──► Consumer: Email
                    │
Publisher ──► Exchange ──── Queue B ──► Consumer: SMS
  (Fanout)          │
                    └── Queue C ──► Consumer: Push

Mỗi queue nhận BẢN COPY của message
Mỗi consumer xử lý theo cách riêng

4. Apache Kafka

4.1 Architecture

┌─────────────────────────────────────────────────────┐
│                    Kafka Cluster                     │
│                                                      │
│  Topic: "orders" (3 partitions, RF=3)                │
│                                                      │
│  ┌─────────┐  ┌─────────┐  ┌─────────┐             │
│  │Broker 1 │  │Broker 2 │  │Broker 3 │             │
│  │         │  │         │  │         │             │
│  │P0(leader│  │P1(leader│  │P2(leader│             │
│  │P1(replica│  │P2(replica│  │P0(replica│             │
│  │P2(replica│  │P0(replica│  │P1(replica│             │
│  └─────────┘  └─────────┘  └─────────┘             │
│                                                      │
│  ZooKeeper / KRaft: Cluster metadata                │
└─────────────────────────────────────────────────────┘

Producer ──► Partition (key-based routing)
Consumer Group ──► Mỗi partition chỉ 1 consumer

4.2 Kafka vs RabbitMQ

FeaturesRabbitMQKafka
ModelMessage QueueEvent Log
Throughput~50K msg/s~1M msg/s
Message retentionConsumed → deletedRetained (configurable)
OrderingPer queuePer partition
ReplayNoYes (offset reset)
Use casesTask distributionEvent streaming, logs
RoutingExchange + bindingTopic + partition
ProtocolsAMQPCustom binaries

5. Delivery Guarantees

5.1 At-Most-Once

Producer → Broker: Send message (fire & forget)
Nếu network error → Message mất
Đơn giản, nhanh, nhưng có thể mất data

Use case: Metrics, logs (mất vài message OK)

5.2 At-Least-Once

Producer → Broker: Send message
Broker → Producer: ACK
Nếu không nhận ACK → Producer gửi lại
→ Message có thể bị duplicate!

Use case: Most applications (xử lý duplicate bằng idempotency)

5.3 Exactly Once

Rất khó đạt được thực sự. Thường là:
  At-Least-Once + Idempotent Consumer

Idempotent Consumer:
  message_id = "order-123-email"
  
  IF NOT EXISTS processed_messages[message_id]:
      process(message)
      INSERT processed_messages[message_id]
  ELSE:
      skip (đã xử lý rồi)

6. Dead Letter Queue (DLQ)

Main Queue ──► Consumer ──► Xử lý thành công ✅
                 │
                 ├── Fail lần 1 → Retry
                 ├── Fail lần 2 → Retry
                 ├── Fail lần 3 → Retry
                 │
                 └── Fail lần 4 → Dead Letter Queue
                                   │
                                   ▼
                              ┌─────────┐
                              │  DLQ    │
                              │ Manual  │
                              │ Review  │
                              └─────────┘

DLQ chứa messages không thể xử lý
Ops team review và quyết định:
  - Fix bug → Replay message
  - Bad data → Discard
  - Dependency down → Replay sau khi fix

7. Backpressure

Vấn đề: Producer nhanh hơn Consumer

 Producer: 10K msg/s ────► Queue ────► Consumer: 2K msg/s
                            │
                       Queue grows!
                       Memory fills!
                       System crash!

Giải pháp:
  1. Rate limiting: Producer giới hạn tốc độ
  2. Buffering: Queue có max size, reject khi đầy
  3. Scaling: Thêm consumers
  4. Sampling: Drop non-critical messages
  5. Flow control: Broker báo Producer slow down

Summary

ToolsTypeBest For
RabbitMQMessage BrokerTask distribution, routing
KafkaEvent StreamingHigh throughput, event log
SQSCloud QueueSimple AWS workloads
Redis StreamsIn-memory streamReal-time, lightweight
NATSMessagingCloud-native, low latency
Celery/SidekiqTask QueueBackground jobs

Exercises

  1. Queue Design: E-commerce checkout includes: payment, inventory, email, SMS, analytics. Design message flow. Which pattern to use: Work Queue or Pub/Sub?

  2. Idempotency: Payment service receives message "charge $100 for order-123". Network timeout, producer sends again. How to make sure not to charge twice?

  3. DLQ Strategy: You detect 500 messages in the DLQ of the email service. 80% due to invalid email address, 20% due to SMTP server timeout. Write a treatment plan.