VibeKoding / Ensiklopedia Β· Fondasi KuatEnsiklopedia Β· Fondasi Kuat / Principles of Message Queues and Event-Driven ArchitecturePrinciples of Message Queues and Event-Driven Architecture
VK

Principles of Message Queues and Event-Driven ArchitecturePrinciples of Message Queues and Event-Driven Architecture

πŸ“š Ensiklopedia Β· Fondasi KuatEnsiklopedia Β· Fondasi Kuat 🌏 Dual Bahasa (ID / EN) ⚑ VibeKoding Native

Ensiklopedia VibeKoding: Principles of Message Queues and Event-Driven Architecture.Ensiklopedia VibeKoding: Principles of Message Queues and Event-Driven Architecture.

πŸ’‘ Tips PraktisπŸ’‘ Pro Tip

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.

------

1. Motivation for the "Message Queues"1. Motivation for the "Message Queues"

1.1 A Real-World Case: The Evolution of Taobao's Order System1.1 A Real-World Case: The Evolution of Taobao's Order System

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
⚠️ Catatan Keamanan / Peringatan⚠️ Warning / Security Note

- 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)
πŸ’‘ Tips PraktisπŸ’‘ Pro Tip

- 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

1.2 Everyday Analogies for Message Queues1.2 Everyday Analogies for Message Queues

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":

------

2. Overview of a Message Queue (Definition + Core Three Elements)2. Overview of a Message Queue (Definition + Core Three Elements)

2.1 Overview of a "Message Queue"2.1 Overview of a "Message Queue"

πŸ’‘ Tips PraktisπŸ’‘ Pro Tip

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).

2.2 The Three Core Elements of a Message Queue2.2 The Three Core Elements of a Message Queue

Element 1: ProducerElement 1: Producer

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).

πŸ“– Konsep PentingπŸ“– Core Concept

- 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

Element 2: ConsumerElement 2: Consumer

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.

πŸ“– Konsep PentingπŸ“– Core Concept

- 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

Element 3: Broker (Message Broker)Element 3: Broker (Message Broker)

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.

πŸ“– Konsep PentingπŸ“– Core Concept

- 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

------

3. Core Problem 1: Approach to Decoupling Systems and Avoid "Pulling One Thread and Moving the Whole System"3. Core Problem 1: Approach to Decoupling Systems and Avoid "Pulling One Thread and Moving the Whole System"

3.1 The Tragedy of Tight Coupling: One Service Goes Down, Everything Falls3.1 The Tragedy of Tight Coupling: One Service Goes Down, Everything Falls

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 β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
πŸ’‘ Tips PraktisπŸ’‘ Pro Tip

| 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 |

3.2 Decoupling Solution: Introduce Message Queue as "Middle Layer"3.2 Decoupling Solution: Introduce Message Queue as "Middle Layer"

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β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

πŸ’‘ Tips PraktisπŸ’‘ Pro Tip

| 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 |

3.3 The Essence of Decoupling: From "Direct Calls" to "Event-Driven"3.3 The Essence of Decoupling: From "Direct Calls" to "Event-Driven"

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

------

4. Core Problem 2: Approach to handling Traffic Spikes with Peak Shaving4. Core Problem 2: Approach to handling Traffic Spikes with Peak Shaving

4.1 Flash Sale Scenario: Approach to handling 100K QPS Smoothly4.1 Flash Sale Scenario: Approach to handling 100K QPS Smoothly

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)
πŸ’‘ Tips PraktisπŸ’‘ Pro Tip

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.

4.2 Peak Shaving Solution: Message Queue as "Reservoir"4.2 Peak Shaving Solution: Message Queue as "Reservoir"

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) β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜ β”‚ β”‚ β”‚ β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

4.3 Mathematical Principles of Peak Shaving4.3 Mathematical Principles of Peak Shaving

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

------

5. Core Problem 3: Approach to ensuring Messages Are Not Lost, Not Duplicated, and In Order5. Core Problem 3: Approach to ensuring Messages Are Not Lost, Not Duplicated, and In Order

5.1 Message Reliability: Three Lines of Defense5.1 Message Reliability: Three Lines of Defense

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.

⚠️ Catatan Keamanan / Peringatan⚠️ Warning / Security Note

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

5.2 Approach to handling Duplicate Message Consumption5.2 Approach to handling Duplicate Message Consumption

Message duplication can occur in the following scenarios:Message duplication can occur in the following scenarios:

  1. Producer retry: Producer sends message but doesn't receive ACK, retries sending the same messageProducer retry: Producer sends message but doesn't receive ACK, retries sending the same message
  2. Consumer ACK timeout: Consumer finishes processing but ACK times out, Broker redeliversConsumer ACK timeout: Consumer finishes processing but ACK times out, Broker redelivers
  3. Network jitter: Consumer ACK doesn't reach Broker, Broker considers message unconsumedNetwork jitter: Consumer ACK doesn't reach Broker, Broker considers message unconsumed
  4. Consumer restart: After consumer restarts, re-consumes the same batch of messagesConsumer restart: After consumer restarts, re-consumes the same batch of messages
  5. πŸ’‘ Tips PraktisπŸ’‘ Pro Tip

    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.

    ------

    6. Practice: Approach to choosing a Message Queue6. Practice: Approach to choosing a Message Queue

    6.1 Comparison of Four Mainstream Message Queues6.1 Comparison of Four Mainstream Message Queues

    FeatureRabbitMQKafkaRocketMQRedis Stream
    PositioningTraditional MQDistributed log streamE-commerce-grade MQLightweight queue
    Throughput~10K/s~1M/s~100K/s~50K/s
    LatencyMicrosecondsMillisecondsMillisecondsMilliseconds
    ReliabilityHigh (persistence)High (multi-replica)High (sync flush)Medium (AOF)
    Message replayNot supportedSupportedSupportedSupported
    Transactional messagesSupported (weak)Not supportedSupported (strong)Not supported
    Delayed messagesSupportedNot supportedSupportedNot supported
    Use casesTraditional enterprise appsLogs, big dataE-commerce, financeSmall-scale apps
    πŸ’‘ Tips PraktisπŸ’‘ Pro Tip

    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) ``

    ------

    7. Summary: Message Queue Design Principles7. Summary: Message Queue Design Principles

    7.1 Core Principles Review7.1 Core Principles Review

    PrincipleMeaningPractice Points
    DecouplingServices don't directly depend on each otherCommunicate via message queue; consumer failure doesn't affect producer
    Peak shavingSmooth traffic fluctuationsMessage queue as reservoir; consumers process at constant rate
    ReliabilityMessages not lostProducer ACK + Broker persistence + Consumer ACK
    IdempotencyDuplicate consumption has no effectBusiness-level idempotency guarantees (unique keys, state machines)
    OrderingMessage order guaranteeSingle-partition ordering or consumer-side sorting

    7.2 Design Checklist7.2 Design Checklist

    Before introducing a message queue, ask yourself:Before introducing a message queue, ask yourself:

    • [ ] Do you really need a message queue? (Simple async can use thread pools)[ ] Do you really need a message queue? (Simple async can use thread pools)
    • [ ] Is message loss acceptable? (Determines reliability level)[ ] Is message loss acceptable? (Determines reliability level)
    • [ ] Will message duplication affect the business? (Determines idempotency investment)[ ] Will message duplication affect the business? (Determines idempotency investment)
    • [ ] Is message order important? (Determines partition strategy)[ ] Is message order important? (Determines partition strategy)
    • [ ] What's the consumer processing capacity? (Determines queue size and alert thresholds)[ ] What's the consumer processing capacity? (Determines queue size and alert thresholds)
    • [ ] How to handle consumption failures? (Determines retry and dead letter strategies)[ ] How to handle consumption failures? (Determines retry and dead letter strategies)

    ------

    8. Glossary8. Glossary

    TermFull NameDescription
    MQMessage QueueMiddleware 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.
    ACKAcknowledgmentConsumer confirms to Broker after processing a message.
    Pub/SubPublish/SubscribeA messaging pattern where one message can be received by multiple consumers.
    P2PPoint-to-PointA messaging pattern where one message can only be received by one consumer.
    DLQDead Letter QueueStores 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.