はじめに
同期処理: ユーザーが「注文」をクリック → 30 秒待ちます (注文の作成、メールの送信、送料の計算、在庫の更新など)。遅すぎる! メッセージ キュー により、非同期処理が可能になり、待ち時間が短縮され、信頼性が向上します。
1. 同期と非同期
1.1 同期の問題
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 メッセージキューによる非同期
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: 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 アーキテクチャ
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 ワークキューのパターン
┌── 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 パブリッシュ/サブスクライブ パターン
┌── 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 アーキテクチャ
┌─────────────────────────────────────────────────────┐
│ 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 と RabbitMQ の比較
| 特長 | ラビットMQ | カフカ |
|---|---|---|
| モデル | メッセージキュー | イベントログ |
| スループット | ~50,000 メッセージ/秒 | ~100万メッセージ/秒 |
| メッセージの保存 | 消費→削除 | 保持 (構成可能) |
| 注文 | キューごと | パーティションごと |
| リプレイ | いいえ | はい (オフセットリセット) |
| 使用例 | タスクの分散 | イベントストリーミング、ログ |
| ルーティング | 交換 + バインディング | トピック + パーティション |
| プロトコル | AMQP | カスタムバイナリ |
5. 配送保証
5.1 最大 1 回
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 少なくとも 1 回
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 1 回だけ
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. デッドレターキュー (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. バックプレッシャー
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
概要
| ツール | タイプ | 最適な用途 |
|---|---|---|
| ラビットMQ | メッセージブローカー | タスクの分散、ルーティング |
| カフカ | イベントストリーミング | 高スループット、イベント ログ |
| SQS | クラウドキュー | シンプルな AWS ワークロード |
| Redis ストリーム | インメモリストリーム | リアルタイム、軽量 |
| ナッツ | メッセージ | クラウドネイティブ、低遅延 |
| セロリ/サイドキック | タスクキュー | バックグラウンドジョブ |
演習
-
キューの設計: 電子商取引のチェックアウトには、支払い、在庫、電子メール、SMS、分析が含まれます。メッセージ フローを設計します。ワークキューまたはパブリッシュ/サブスクライブのどちらのパターンを使用するか?
-
冪等性: 支払いサービスは「order-123 に $100 を請求します」というメッセージを受け取ります。ネットワークがタイムアウトしたため、プロデューサーが再度送信しました。二重充電をしないようにするにはどうすればよいですか?
-
DLQ 戦略: 電子メール サービスの DLQ で 500 のメッセージを検出しました。 80% は無効な電子メール アドレスが原因で、20% は SMTP サーバーのタイムアウトが原因です。治療計画を書きます。