Ensiklopedia VibeKoding: Principles of Message Queues and Event-Driven Architecture.Ensiklopedia VibeKoding: Principles of Message Queues and Event-Driven Architecture.
When systems are tightly coupled and traffic spikes, how do you ensure the critical path remains stable? Message queues are the "buffer" and "decoupler" of modern distributed systems. This article uses real-world cases (restaurant queuing, package sorting, flash sale systems) to deeply understand the design philosophy and engineering practices of message queues.When systems are tightly coupled and traffic spikes, how do you ensure the critical path remains stable? Message queues are the "buffer" and "decoupler" of modern distributed systems. This article uses real-world cases (restaurant queuing, package sorting, flash sale systems) to deeply understand the design philosophy and engineering practices of message queues.
------
In 2012, Taobao's order system suffered a severe outage. At midnight on Double 11, traffic flooded in instantly. The order service directly called the inventory service, payment service, logistics service... the entire chain collapsed like dominoes.In 2012, Taobao's order system suffered a severe outage. At midnight on Double 11, traffic flooded in instantly. The order service directly called the inventory service, payment service, logistics service... the entire chain collapsed like dominoes.
The architecture at the time (tight coupling):The architecture at the time (tight coupling):
CODE User places order β Order service β Sync call inventory service β Sync call payment service β Sync call logistics service β β β Response 200ms Response 500ms Response 300ms
- Total response time = 200 + 500 + 300 = 1000ms (user waits 1 second) - Inventory service down β Order service also goes down (thread pool exhausted) - Payment service slows down β Entire chain is dragged down - Cannot scale horizontally β Can only scale vertically (expensive and limited)- Total response time = 200 + 500 + 300 = 1000ms (user waits 1 second) - Inventory service down β Order service also goes down (thread pool exhausted) - Payment service slows down β Entire chain is dragged down - Cannot scale horizontally β Can only scale vertically (expensive and limited)
Improved architecture (introducing message queues):Improved architecture (introducing message queues):
CODE User places order β Order service β Send "order created" message β Immediate return (50ms) β Message queue (Kafka) β βββββββββββββββ¬ββββββββββββββ¬ββββββββββββββ βΌ βΌ βΌ βΌ Inventory Payment Logistics Notification service service service service (async deduct) (async process) (async create) (async send)
- User response time = 50ms (20x experience improvement) - Inventory service down β Messages stay in queue, continue processing after recovery - Payment service slows down β Doesn't affect order creation - Can scale horizontally β Just add more consumer instances- User response time = 50ms (20x experience improvement) - Inventory service down β Messages stay in queue, continue processing after recovery - Payment service slows down β Doesn't affect order creation - Can scale horizontally β Just add more consumer instances
Restaurant Queuing SystemRestaurant Queuing System
Imagine going to a popular restaurant:Imagine going to a popular restaurant:
A message queue is the software system's "queuing system":A message queue is the software system's "queuing system":
------
Message Queue (MQ) is a container for storing messages. Producers put messages in, consumers take messages out for processing. It enables "asynchronous communication" β the sender doesn't need to wait for the receiver to finish processing. Synchronous vs Asynchronous: - Synchronous: Like a phone call β the other party must answer to communicate - Asynchronous: Like sending a text β you send it, they read it when available It's like calling a friend (synchronous) vs sending them a message (asynchronous).Message Queue (MQ) is a container for storing messages. Producers put messages in, consumers take messages out for processing. It enables "asynchronous communication" β the sender doesn't need to wait for the receiver to finish processing. Synchronous vs Asynchronous: - Synchronous: Like a phone call β the other party must answer to communicate - Asynchronous: Like sending a text β you send it, they read it when available It's like calling a friend (synchronous) vs sending them a message (asynchronous).
Responsibility: Create and send messages to the queue.Responsibility: Create and send messages to the queue.
Analogy: The producer is like a "sender," delivering letters (messages) to the post office (queue).Analogy: The producer is like a "sender," delivering letters (messages) to the post office (queue).
- Send method: Synchronous send (reliable but blocking) vs asynchronous send (high performance but needs callback handling) - Message confirmation: Wait for Broker confirmation (At Least Once) vs fire-and-forget (At Most Once) - Failure handling: Retry strategy, local log backup, dead letter queue- Send method: Synchronous send (reliable but blocking) vs asynchronous send (high performance but needs callback handling) - Message confirmation: Wait for Broker confirmation (At Least Once) vs fire-and-forget (At Most Once) - Failure handling: Retry strategy, local log backup, dead letter queue
Responsibility: Get messages from the queue and process them.Responsibility: Get messages from the queue and process them.
Analogy: The consumer is like a "recipient," taking letters (messages) from the mailbox (queue) and processing them.Analogy: The consumer is like a "recipient," taking letters (messages) from the mailbox (queue) and processing them.
- Consumption mode: Push mode (Broker actively pushes) vs Pull mode (consumer actively pulls) - Consumption confirmation: Auto ACK (efficient but may lose messages) vs manual ACK (reliable but needs timeout handling) - Concurrency control: Single-threaded sequential consumption vs multi-threaded parallel consumption - Failure handling: Retry strategy, dead letter queue, compensation mechanism- Consumption mode: Push mode (Broker actively pushes) vs Pull mode (consumer actively pulls) - Consumption confirmation: Auto ACK (efficient but may lose messages) vs manual ACK (reliable but needs timeout handling) - Concurrency control: Single-threaded sequential consumption vs multi-threaded parallel consumption - Failure handling: Retry strategy, dead letter queue, compensation mechanism
Responsibility: Receive, store, and forward messages.Responsibility: Receive, store, and forward messages.
Analogy: The Broker is like a "post office" or "package sorting station," responsible for receiving, sorting, and delivering letters.Analogy: The Broker is like a "post office" or "package sorting station," responsible for receiving, sorting, and delivering letters.
- Storage model: In-memory storage (low latency) vs disk storage (high reliability) - Replication strategy: Primary-secondary replication, multi-replica synchronization - High availability: Cluster deployment, automatic failover - Scalability: Partitions, Sharding- Storage model: In-memory storage (low latency) vs disk storage (high reliability) - Replication strategy: Primary-secondary replication, multi-replica synchronization - High availability: Cluster deployment, automatic failover - Scalability: Partitions, Sharding
------
Scenario recreation: An e-commerce platform's early architectureScenario recreation: An e-commerce platform's early architecture
CODE Order service directly calls downstream services: βββββββββββββββ β Order β β Service β ββββββββ¬βββββββ β βββββββββββββ¬ββββββββββββ¬ββββββββββββ βΌ βΌ βΌ βΌ ββββββββββββ ββββββββββββ ββββββββββββ ββββββββββββ βInventory β βPayment β βLogistics β βSMS β βService β βService β βService β βService β β 200ms β β 500ms β β 300ms β β 100ms β ββββββββββββ ββββββββββββ ββββββββββββ ββββββββββββ
| Pain Point | Specific Manifestation | Consequence | |------|----------|------| | Cascading failure | Inventory service goes down, order service sync call times out | Order service thread pool exhausted, cannot process new requests | | Response latency | Must wait for all downstream service responses | User waits over 1 second, terrible experience | | Difficult to extend | Adding a points service requires modifying order service code | Longer release cycles, increased risk | | Resource waste | Order service must wait for SMS service | Database connections occupied for long periods || Pain Point | Specific Manifestation | Consequence | |------|----------|------| | Cascading failure | Inventory service goes down, order service sync call times out | Order service thread pool exhausted, cannot process new requests | | Response latency | Must wait for all downstream service responses | User waits over 1 second, terrible experience | | Difficult to extend | Adding a points service requires modifying order service code | Longer release cycles, increased risk | | Resource waste | Order service must wait for SMS service | Database connections occupied for long periods |
Architecture after decoupling:Architecture after decoupling:
CODE Order service only sends messages, doesn't care who consumes: βββββββββββββββ β Order β ββSend "order created" messageβββ β Service β β βββββββββββββββ βΌ βββββββββββββββββββββ β Message Queue β β (Kafka/RabbitMQ) β β - Reliable store β β - Multi-replica β β - Order guaranteeβ βββββββββββ¬ββββββββββ β βββββββββββββββββββββββββΌββββββββββββββββββββββββ β β β βΌ βΌ βΌ ββββββββββββββββ ββββββββββββββββ ββββββββββββββββ β Inventory β β Payment β β Logistics β β Service β β Service β β Service β β Subscribe β β Subscribe β β Subscribe β β order eventsβ β order eventsβ β order eventsβ ββββββββββββββββ ββββββββββββββββ ββββββββββββββββ
| Dimension | Before Decoupling | After Decoupling | |------|--------|--------| | Fault isolation | Inventory down = Order down | Inventory down, messages stay in queue, consumed after recovery | | Response time | 1000ms (synchronous wait) | 50ms (return after sending message) | | Extensibility | New service requires changing order code | New service just subscribes to topic | | System complexity | Order service tightly depends on downstream | Order service only depends on message queue || Dimension | Before Decoupling | After Decoupling | |------|--------|--------| | Fault isolation | Inventory down = Order down | Inventory down, messages stay in queue, consumed after recovery | | Response time | 1000ms (synchronous wait) | 50ms (return after sending message) | | Extensibility | New service requires changing order code | New service just subscribes to topic | | System complexity | Order service tightly depends on downstream | Order service only depends on message queue |
Paradigm shift:Paradigm shift:
CODE Traditional thinking (imperative): "Order service commands inventory service: Deduct inventory for me!" β Direct call β High coupling, callee must be online β Caller needs to know callee's interface Event-driven thinking (declarative): "Order service declares: Order has been created. Whoever cares, handle it." β Send event to message queue β Decoupled, consumers can be offline β Producer doesn't need to know consumers exist
------
Scenario recreation: An e-commerce platform's Double 11 flash sale, expected peak 100K QPS, but the database can only handle 1,000 QPS.Scenario recreation: An e-commerce platform's Double 11 flash sale, expected peak 100K QPS, but the database can only handle 1,000 QPS.
Consequences of direct impact:Consequences of direct impact:
CODE User requests βββ App server βββ Database 100K/s 100K/s 1K/s (limit) β Connection pool exhausted Response timeout Database crash β Cascading failure (all services depending on DB go down)
QPS (Queries Per Second): Queries per second, a metric for measuring system throughput. 100K QPS means 100,000 requests per second, like 100,000 people rushing into a store simultaneously.QPS (Queries Per Second): Queries per second, a metric for measuring system throughput. 100K QPS means 100,000 requests per second, like 100,000 people rushing into a store simultaneously.
Architecture design:Architecture design:
CODE βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β Flash Sale System Architecture β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ€ β β β Layer 1: Gateway (hard rate limiting) β β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β β β - Token bucket: 100K/s β 10K/s (drop 90% of requests) β β β β - CDN caches static resources (product detail pages) β β β β - CAPTCHA / queue page (first layer of peak shaving) β β β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β β β β β βΌ β β Layer 2: Service (soft rate limiting) β β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β β β - Nginx rate limiting: 10K/s β 5K/s β β β β - Redis pre-deduct inventory (atomic operation): β β β β * Use Lua script for atomicity β β β β * Insufficient stock β return "Sold out" directly β β β β - Generate order token (queue voucher) β β β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β β β β β βΌ β β Layer 3: Message Queue (core peak shaving) β β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β β β Kafka/RocketMQ: β β β β - Batch write: 5K/s β 1K/s (matching DB capacity) β β β β - Message persistence: disk write guarantees no message loss β β β β - Multi-partition parallel consumption: boost throughput β β β β - Consumer offset management: support failure recovery β β β β β β β β Key metrics monitoring: β β β β - Produce Rate β β β β - Consume Rate β β β β - Lag (message backlog) β β β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β β β β β βΌ β β Layer 4: Consumer (async processing) β β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β β β Order processing consumers (multiple instances): β β β β - Pull messages from Kafka (1K/s, matching DB capacity) β β β β - DB transaction: create order + deduct inventory β β β β - Update order status to "Created" β β β β - Send order creation success notification (email/SMS/push) β β β β - Confirm message consumption (ACK) β β β β β β β β Consumer scaling strategy: β β β β - When Lag > 10,000, auto-scale up consumer instances β β β β - When Lag < 1,000, scale down consumer instances (save cost) β β β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ β β β βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
Traffic smoothing effect:Traffic smoothing effect:
CODE Original traffic (spike): Smoothed traffic: 100K/s β β±β² 1K/s βββββββββββββββββ β β± β² β β β± β² β 1K/sββ± β² 0/s β ββββββββββββββββ βββββββββββββββββ 0s 1s 2s 0s 20s Original: 100K/s peak, lasting 1 second Smoothed: 1K/s constant rate, lasting 100 seconds
Key formulas:Key formulas:
CODE Queue length = Producer rate Γ Duration - Consumer rate Γ Duration = 100,000 Γ 1 - 1,000 Γ 1 = 99,000 messages (peak queue backlog) Time to consume all messages = Queue length / Consumer rate = 99,000 / 1,000 = 99 seconds
------
Messages can be lost at three stages: during producer sending, during Broker storage, and during consumer processing.Messages can be lost at three stages: during producer sending, during Broker storage, and during consumer processing.
Defense 1: Producer ACK - When sending a message, wait for the Broker to confirm receipt - If no confirmation received, retry or log locally Defense 2: Broker Persistence - Write messages to disk, not just in memory - Multi-replica synchronization to ensure no data loss Defense 3: Consumer ACK - After processing a message, manually confirm (ACK) - If processing fails, don't confirm; Broker will redeliverDefense 1: Producer ACK - When sending a message, wait for the Broker to confirm receipt - If no confirmation received, retry or log locally Defense 2: Broker Persistence - Write messages to disk, not just in memory - Multi-replica synchronization to ensure no data loss Defense 3: Consumer ACK - After processing a message, manually confirm (ACK) - If processing fails, don't confirm; Broker will redeliver
Message duplication can occur in the following scenarios:Message duplication can occur in the following scenarios:
Idempotency: Executing the same operation multiple times produces the same result as executing it once. Everyday idempotency examples: - Idempotent: Pressing an elevator button (pressing 10 times or once, the elevator still comes) - Non-idempotent: Bank transfer (transferring $10, executing twice transfers $20) Technical solution: Generate a unique ID for each message; check if already processed before handling.Idempotency: Executing the same operation multiple times produces the same result as executing it once. Everyday idempotency examples: - Idempotent: Pressing an elevator button (pressing 10 times or once, the elevator still comes) - Non-idempotent: Bank transfer (transferring $10, executing twice transfers $20) Technical solution: Generate a unique ID for each message; check if already processed before handling.
------
| Feature | RabbitMQ | Kafka | RocketMQ | Redis Stream |
|---|---|---|---|---|
| Positioning | Traditional MQ | Distributed log stream | E-commerce-grade MQ | Lightweight queue |
| Throughput | ~10K/s | ~1M/s | ~100K/s | ~50K/s |
| Latency | Microseconds | Milliseconds | Milliseconds | Milliseconds |
| Reliability | High (persistence) | High (multi-replica) | High (sync flush) | Medium (AOF) |
| Message replay | Not supported | Supported | Supported | Supported |
| Transactional messages | Supported (weak) | Not supported | Supported (strong) | Not supported |
| Delayed messages | Supported | Not supported | Supported | Not supported |
| Use cases | Traditional enterprise apps | Logs, big data | E-commerce, finance | Small-scale apps |
Decision tree: `` Choosing a message queue: β ββ Need transactional messages (distributed transactions)? β ββ Yes β RocketMQ (first choice) or RabbitMQ β ββ No β continue β ββ Need to process massive logs/real-time streams? β ββ Yes β Kafka (first choice) β ββ No β continue β ββ QPS > 10K/s? β ββ Yes β RocketMQ or Kafka β ββ No β continue β ββ Need complex routing (e.g., header matching)? β ββ Yes β RabbitMQ β ββ No β continue β ββ Already have Redis infrastructure? β ββ Yes β Redis Stream (quick start) β ββ No β RabbitMQ (full-featured, moderate learning curve) ``Decision tree: `` Choosing a message queue: β ββ Need transactional messages (distributed transactions)? β ββ Yes β RocketMQ (first choice) or RabbitMQ β ββ No β continue β ββ Need to process massive logs/real-time streams? β ββ Yes β Kafka (first choice) β ββ No β continue β ββ QPS > 10K/s? β ββ Yes β RocketMQ or Kafka β ββ No β continue β ββ Need complex routing (e.g., header matching)? β ββ Yes β RabbitMQ β ββ No β continue β ββ Already have Redis infrastructure? β ββ Yes β Redis Stream (quick start) β ββ No β RabbitMQ (full-featured, moderate learning curve) ``
------
| Principle | Meaning | Practice Points |
|---|---|---|
| Decoupling | Services don't directly depend on each other | Communicate via message queue; consumer failure doesn't affect producer |
| Peak shaving | Smooth traffic fluctuations | Message queue as reservoir; consumers process at constant rate |
| Reliability | Messages not lost | Producer ACK + Broker persistence + Consumer ACK |
| Idempotency | Duplicate consumption has no effect | Business-level idempotency guarantees (unique keys, state machines) |
| Ordering | Message order guarantee | Single-partition ordering or consumer-side sorting |
Before introducing a message queue, ask yourself:Before introducing a message queue, ask yourself:
------
| Term | Full Name | Description |
|---|---|---|
| MQ | Message Queue | Middleware for asynchronous communication, decoupling producers and consumers. |
| Producer | - | The party that sends messages. |
| Consumer | - | The party that receives and processes messages. |
| Broker | - | The server program that stores and forwards messages. |
| Topic | - | Logical categorization of messages (e.g., "orders"). |
| Queue | - | Physical container storing messages. |
| Partition | - | A Kafka concept; one Topic can be split into multiple Partitions for higher concurrency. |
| ACK | Acknowledgment | Consumer confirms to Broker after processing a message. |
| Pub/Sub | Publish/Subscribe | A messaging pattern where one message can be received by multiple consumers. |
| P2P | Point-to-Point | A messaging pattern where one message can only be received by one consumer. |
| DLQ | Dead Letter Queue | Stores messages that cannot be consumed. |
| Idempotence | - | Multiple executions produce the same result. |
| Throughput | - | Number of messages processed per unit time. |
| Latency | - | Time difference from message send to receipt. |
| Persistence | - | Messages written to disk, not just stored in memory. |
| Replication | - | Messages copied to multiple nodes for high availability. |
| Transaction Message | - | Guarantees consistency between local transaction and message sending. |
| Backpressure | - | When consumers can't keep up, they notify producers to slow down. |
| Offset | - | The consumer's consumption position within a partition. |
| Rebalance | - | Reassigning partitions when consumer group members change. |