MSG
QUEUES
CONSUMER GROUPS · PARTITIONS · EXACTLY-ONCE
// WITHOUT message queue: tight coupling, brittle [Order Service] ──sync──→ [Inventory Service] // What if inventory is down? [Order Service] ──sync──→ [Email Service] // What if email is slow (2s)? [Order Service] ──sync──→ [Analytics Service] // User waits for all 3! // WITH message queue: decoupled, resilient, fast [Order Service] ──→ [Topic: order-placed] ──→ [Inventory] ← independent ──→ [Email] ← independent ──→ [Analytics] ← independent Order Service returns in <1ms. Downstream services process asynchronously. If email is down → messages queue up → processed when it recovers. Each consumer processes at its own rate. No cascading failures.
Consumer A2 → P1
Consumer A3 → P2
3 consumers = 3 partitions ✓
1 consumer handles all — slower, but independent
Consumer C2 → P1
Consumer C3 → P2
fully parallel ✓
// replication.factor=3 → 1 leader + 2 ISR replicas per partition // ISR = In-Sync Replicas (have replicated all leader messages) acks=0: Producer doesn't wait for ACK. Fastest, no durability guarantee. acks=1: Leader ACKs after writing. Fast, but replica may not have it yet. acks=all:All ISR ACKs before producer gets confirmation. Strongest guarantee. // With acks=all + min.insync.replicas=2 + replication.factor=3: // → Can lose 1 broker with ZERO data loss // → Brokers 1 (leader) + Broker 2 (replica) both have message before ACK // → Broker 1 dies → Broker 2 becomes leader → no data lost // Partition key routing: producer.send("order-placed", userId, orderEvent); // hash(userId) % numPartitions → same userId → same partition → ordered
Failure: consumer crashes after commit, before processing → message LOST.
(loss is acceptable)
Failure: consumer crashes after processing, before commit → message DUPLICATED.
idempotent consumer pattern
No loss, no duplicates. Requires enable.idempotence=true + transactional.id.
inventory — real harm from duplication
public void processPayment(PaymentEvent e) { // Check idempotency key — has this message been processed before? if (db.exists("processed:" + e.idempotencyKey)) { log.info("Duplicate — skipping: {}", e.idempotencyKey); return; // Silently no-op on duplicate }<span class="mb4-cm">// Atomic: process + mark as processed in same DB transaction</span> db.<span class="mb4-fn">transaction</span>(() -> { db.<span class="mb4-fn">debitAccount</span>(e.accountId, e.amount); db.<span class="mb4-fn">markProcessed</span>(<span class="mb4-str">"processed:"</span> + e.idempotencyKey); }); <span class="mb4-cm">// Now commit Kafka offset — at-least-once is effectively exactly-once</span> consumer.<span class="mb4-fn">commitSync</span>();}
// Key: idempotencyKey must uniquely identify the business operation // Options: UUID in message, (userId + orderId + action), event sequence number
| ASPECT | KAFKA | RABBITMQ |
|---|---|---|
| Delivery model | Pull (consumers fetch at own pace) | Push (broker delivers to consumer) |
| Message retention | Durable log — days to years after delivery | Deleted on ACK — ephemeral |
| Throughput | Millions of messages/second | Hundreds of thousands/sec |
| Ordering guarantee | Within a partition (by key) | Within a queue (single consumer) |
| Replay history | Yes — seek to any offset | No — deleted after consume |
| Multiple consumers | N independent consumer groups | Competing consumers (one gets each msg) |
| Routing logic | Topic/partition key only | Exchanges: direct, topic, fanout, headers |
| Message TTL | Topic-level retention only | Per-message TTL, priority queues |
✓ User activity stream (audit log needed)
✓ CDC: DB changes → Elasticsearch
✓ Real-time analytics pipeline
✓ Microservice event backbone
✓ New service needs historical data (replay)
✓ Task distribution to N worker processes
✓ Routing by message type to different queues
✓ Per-job TTL (expire unprocessed jobs)
✓ Priority queue (high-priority tasks first)
✓ Simple job scheduler without replay needs
Group inventory-svc, Group email-svc,
Group analytics-svc, Group fraud-svc
→ send to topic.dlq
→ alert on-call → fix → replay
Action: scale consumer group
Ceiling: max_consumers = num_partitions
→ hash(userId) % numPartitions
→ same partition = in-order
After: [u2:v1][u1:v3] (latest per key)
Use: user profiles, config, inventory
INSERT INTO orders ...
INSERT INTO outbox (event, payload)
COMMIT → CDC picks up → Kafka
// Example: order event stream // 1M orders/day, peak 50× averagePeak events/sec = 1M orders/day ÷ 86,400 × 50 = ~580 events/sec Event size = 1 KB Peak throughput = 580 × 1KB = ~0.6 MB/sec // Single partition max throughput: ~100 MB/sec write Partitions needed = 0.6 MB/sec ÷ 100 MB/sec = 1 partition (use 12 for growth headroom) // Storage (7-day retention, 3 replicas): Daily = 580 events/sec × 86,400 × 1 KB = ~50 GB/day Total = 50 GB × 7 days × 3 replicas = ~1.05 TB // General rules: // num_partitions ≥ max_consumers_in_any_group // num_partitions = target_throughput_MB_s ÷ throughput_per_partition_MB_s // Start with 12–24, easier to add partitions than subtract
For each, choose at-most-once / at-least-once / exactly-once. State the failure scenario and idempotency strategy:
- Real-time page view counter for analytics dashboard
- Bank transfer between two accounts triggered by Kafka event
- Email notification: "Your order has shipped" (user receives one email)
- Inventory decrement when an order is placed (oversell = bad)
- User activity feed update — showing what friends liked
For each: what is the exact failure scenario if you choose wrong?
For each: choose the partition key, state the ordering guarantee provided, and identify any potential hotkey risk:
- E-commerce order events: created → paid → shipped → delivered (must be in order)
- WhatsApp group chat messages (order within a conversation matters)
- Real-time stock prices for 10,000 tickers (Apple trades 1000× more than a small cap)
- IoT sensor readings from 100K devices
- User login/logout events (session coherence required)
- 8 microservices all need to react to every new user signup — each does something different
- Background job system that sends weekly digest emails to 10M users
- Fraud detection pipeline that must audit every financial transaction for 5 years
- Real-time bidding system where each ad impression must be handled by exactly ONE bidder
- CDC pipeline streaming database row changes to Elasticsearch for search indexing
Context: 5M active drivers sending GPS every 4 seconds. Ride lifecycle events (requested, matched, started, completed, rated). Surge pricing recalculated per zone per minute. Real-time analytics + historical audit.
Design complete Kafka architecture. For each topic:
- Topic name and purpose
- Partition key choice and ordering guarantee
- Delivery semantic (with justification)
- Retention policy (with justification)
- Which consumer groups consume it and what they do
Calculate: peak events/sec total, storage/day total, minimum partitions needed.
Short code generation (base62, MD5) · Redirect latency <10ms · Hot URL caching
Analytics pipeline · Rate limiting · Custom aliases · TTL expiry