Sizing Consumer Worker Fleets to Drain Queue Backlogs
When an upstream service experiences a traffic surge or downstream outages cause message queues (like Apache Kafka, RabbitMQ, or AWS SQS) to accumulate hundreds of thousands of unprocessed events, engineers must scale consumer capacity to drain the backlog while concurrently absorbing live incoming production traffic.
Key Considerations: Kafka Partitions vs Worker Threads
- Kafka Partition Bottleneck: In Apache Kafka, each partition can only be read by at most one consumer thread in a consumer group. If you have 8 partitions, scaling to 32 worker pods will leave 24 pods idle.
- Downstream Database Shock: Aggressively scaling consumers to drain a backlog in 5 minutes can generate tens of thousands of database write queries per second, knocking over PostgreSQL, Redis, or third-party APIs.