Shadow Consumer Group Pattern Clears 1.2M Kafka Backlog in 11 Minutes

A real-world case study demonstrates a six-step SOP using shadow consumer groups and fallback consumers to resolve a 1.2 million message Kafka backlog in 11 minutes while maintaining data consistency and preventing downstream database overload.

Tech Freedom Circle
Tech Freedom Circle
Tech Freedom Circle
Shadow Consumer Group Pattern Clears 1.2M Kafka Backlog in 11 Minutes

Three Real-World Backlog Incidents

Case 1: E-commerce Kafka 1.2M Order Backlog (Shadow Group Scenario)

Business: Order topic, 3M+ daily orders, 6 partitions, consumer group main-consumer-group, 3 pods, concurrency=3 (effective concurrency 6 due to partition limit).

Symptom: Lag spiked to 1.2M+, consumption TPS dropped from 1,800 to 200. Scaling pods to 8 did not help — extra instances stayed idle because partitions capped at 6.

Root cause: New logistics sync logic added synchronous HTTP calls to a third-party logistics API. Nighttime instability caused timeouts; single-message RT jumped from 5 ms to 200+ ms, blocking partition consumption while producers kept pushing.

Case 2: Food Delivery RocketMQ 1M Backlog (Queue Bottleneck)

Business: Flash-sale order flow, 4 queues, 2 consumer nodes normally.

Symptom: Peak traffic pushed backlog past 1M; inventory deduction delayed 10+ minutes, oversell risk. Scaling to 10 pods failed — only 4 queues, so extra consumers fetched nothing.

Root cause: Unoptimized inventory SQL caused slow queries, database latency blocked consumer threads; production rate exceeded consumption rate.

Note: RocketMQ provides native retry queues and DLQ; massive backlog can be replayed via MQAdmin offset reset + temporary consumer group — no custom shadow group needed.

Case 3: Overseas E-commerce Kafka 360K Secondary Incident

After flash sale, 360K messages backlogged, 6-hour delay. Developers restarted main consumer at full speed without throttling or degradation, instantly driving downstream DB TPS to 1,000+, exhausting connection pool and CPU, causing 3-hour total outage.

Root cause: Massive replay without traffic shaping or rate limiting turned catch-up into a cascading failure.

Six-Step SOP for 1.2M Kafka Backlog

Overall SOP: Observe metrics → Emergency rate limiting → Dirty message isolation + DLQ → Deploy shadow group to replay backlog (write to repair table) → Fix main consumer with batch + idempotency → Standby fallback group → Data consistency verification → Shadow group cleanup & offset wrap-up

Step 1: Emergency Rate Limiting (Stop the Bleeding)

Monitoring showed producer TPS stable at 1,600/s, consumer TPS only 200, Lag rising ~80K/min.

Applied Sentinel gateway rule on order-produce-api to cap global produce QPS at 1,300, allowing core orders through while blocking non-marketing/test orders.

Distributed tracing pinpointed the block: third-party logistics HTTP call lacked timeout and circuit breaker, stalling consumer threads.

Stop condition: Lag stops rising and plateaus.

// Sentinel gateway rate limit, control total produce QPS for order topic
FlowRule rule = new FlowRule("order-produce-api");
rule.setCount(1300);
rule.setGrade(RuleConstant.FLOW_GRADE_QPS);
FlowRuleManager.loadRules(Lists.newArrayList(rule));

Step 2: Traffic Isolation + Dead Letter Queue (Isolate Bad Messages)

Reused pre-defined DLQ topic order_dlq_topic (3 partitions).

Modified consumer logic: catch consumption/timeout exceptions, route to DLQ instead of endless retry; on DLQ send failure, persist locally + alert to prevent silent loss.

Deployed dedicated DLQ consumer pod (1 instance) to process poison messages without occupying main 6-partition threads; monitor DLQ lag.

Result: Exception messages no longer block partitions; main consumer TPS recovered to 450.

try { handleOrder(record.value()); }
catch (Exception e) {
    kafkaTemplate.send(KafkaConfig.ORDER_DLQ_TOPIC, record.key(), record.value())
        .addCallback(()->{}, ex-> localSaveAndAlert(record));
}
ack.acknowledge();

Step 3: Deploy Shadow Consumer Group (Core Catch-Up Action)

Principle: New group shadow-consumer-group uses independent offset cursor, does not affect main group offsets or real-time order business .

Shadow group does not write to order main table ; only writes to order_backlog_repair repair table, skips logistics push. Separate compensation job later reads repair table and completes logistics push, avoiding dual-write state conflicts.

Configuration: max.poll.records=50 batch fetch, concurrency=5, 3 pods. Partitions = 6 → max effective concurrency 6.

Reset shadow group offset to earliest to replay historical 1.2M messages.

Business degradation: skip third-party logistics remote call, only persist basic order info to repair table.

Downstream DB rate limit: single-instance DB QPS capped at 400, 3 pods combined peak ≤ 1,800 (learning from Case 3).

Observed: Shadow group stable 2,200 TPS consuming history into repair table; main group continues real-time at 1,300 TPS writing to main table.

Total clear time: ~11 minutes (1.2M backlog + ongoing real-time flow).

@KafkaListener(groupId = "shadow-consumer-group")
public void consumeShadow(List<ConsumerRecord<String,String>> records, Acknowledgment ack){
    records.forEach(r -> backlogRepairService.saveToRepairTable(r.value()));
    ack.acknowledge();
}

Step 4: Fix Main Consumer — Batch Consumption + Idempotency

Code fix: logistics API gets 100 ms timeout + circuit breaker; timeout skips push, defers to async compensation.

Enable batch consumption in container factory: max.poll.records=50, keep concurrency=3, 3 pods, effective concurrency 6.

Idempotency two layers: (1) in-batch in-memory HashSet for intra-batch dedup; (2) global idempotency via Redis marker + order DB unique status index — survives restarts.

Result: single-message RT back to 5 ms, main consumer TPS restored to 1,800, real-time lag stays at 0.

Boundary: batch processing failure on a single record routes only that record to DLQ; rest commit offset normally — no full-batch rollback.

for(var r:records){
    if(idempotentService.checkAndMark(r.key())) handleOrder(r.value());
}
ack.acknowledge();

Step 5: Fallback Consumer Group on Standby (Last Line of Defense)

Deploy independent fallback-consumer-group, 3 pods, normally paused, offset static.

Activation switch via Nacos dynamic config: alert when main group lag > 100K for 5 min; ops can one-click enable. Can also integrate KEDA auto-scaling.

Fallback logic writes to separate compensation topic, not core order table, preventing data conflicts.

Verification: fallback group does not compete for partitions with main/shadow — offsets fully isolated.

@KafkaListener(groupId = "fallback-consumer-group")
public void consumeFallback(List<ConsumerRecord<String,String>> records, Acknowledgment ack){
    records.forEach(r -> fallbackProducer.sendCompensateMsg(r.value()));
    ack.acknowledge();
}

Step 6: Post-Clear Consistency Check + Resource Cleanup

Watch shadow group lag; when 0, stop shadow service.

Consistency check: compare order_backlog_repair IDs with main order table; compensation job serially replays logistics push; final audit of order states and total message count confirms no loss, no duplicate orders.

Shadow offset reset to

latest</sub>; retain group config for a while for forensic replay, then delete.

Scale down shadow pods, release CPU/memory.

Purge DLQ poison messages, export logs for root-cause review; keep monitoring DLQ.

Results & Metrics Closure

Initial backlog 1.2M → 0, total clear time ~11 minutes (including continuous real-time inflow).

Shadow peak catch-up TPS: 2,200.

Main real-time TPS stable 1,800, real-time lag 0, zero new user complaints.

Downstream DB peak QPS held ≤ 1,800, avoiding Case 3 style secondary crash.

Zero message loss, zero order state overwrite conflicts; backlog orders completed logistics push via compensation job, final data consistent.

Post-Mortem & Interview Talking Points

Kafka concurrency hard limit = topic partition count; blind pod scaling is useless. Partition expansion possible online but triggers rebalance & cluster jitter — not preferred in emergency.

For consumption-blocking million-scale backlog, shadow group must not write core business table . Write to repair table only; separate compensation job aligns business state, avoiding main/shadow dual-write races.

Shadow replay requires business degradation + downstream rate limiting, otherwise replay floods DB and causes secondary avalanche.

Solution is a combo: gateway rate limit + DLQ isolation + shadow replay + batch consumption + fallback standby. No single tactic handles million-level backlog safely.

Idempotency in two layers: in-batch memory dedup reduces redundant work; global idempotency via Redis/DB unique index survives restarts.

Batch consumption: single record failure → route that record to DLQ, others commit offset normally — no full-batch stall.

Original Source

Signed-in readers can open the original source through BestHub's protected redirect.

Sign in to view source
Republication Notice

This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactadmin@besthub.devand we will review it promptly.

Kafkaincident responsemessage queueidempotencydead letter queuebacklog handlingfallback consumershadow consumer group
Tech Freedom Circle
Written by

Tech Freedom Circle

Crazy Maker Circle (Tech Freedom Architecture Circle): a community of tech enthusiasts, experts, and high‑performance fans. Many top‑level masters, architects, and hobbyists have achieved tech freedom; another wave of go‑getters are hustling hard toward tech freedom.

0 followers
Reader feedback

How this landed with the community

Sign in to like

Rate this article

Was this worth your time?

Sign in to rate
Discussion

0 Comments

Thoughtful readers leave field notes, pushback, and hard-won operational detail here.