# The Trillion Message Kafka Setup at Walmart

## Metadata
- Author: [[ByteByteGo]]
- Full Title: The Trillion Message Kafka Setup at Walmart
- Category: #articles
- Summary: Walmart processes trillions of Kafka messages daily with a setup that ensures 99.99% availability. To manage traffic spikes and reduce costs, they use a system called MPS, which decouples message consumption from Kafka’s partition limits. MPS allows consumer applications to scale independently while maintaining message order and handling failures efficiently.
- URL: https://blog.bytebytego.com/p/the-trillion-message-kafka-setup
## Highlights
- 1 - Consumer Rebalancing
One of the most frequent problems was related to consumer rebalancing.
But what triggers consumer rebalancing in Kafka?
This can happen due to the changing number of consumer instances within a consumer group. ([View Highlight](https://read.readwise.io/read/01m3nkv67wbza4p8ma5zjs3c0g))
-  ([View Highlight](https://read.readwise.io/read/01m3nkv8y4xzcjd5ka1qyac9yt))
- Several scenarios are possible such as:
• A consumer pod may enter or leave a consumer group. This can happen due to Kubernetes deployments, rolling restarts, or automatic scale-ins or scale-outs. Whenever it happens, Kafka needs to redistribute the partitions among the consumers.
• The Kafka broker may believe that a consumer has failed. If the broker has not received a heartbeat from a consumer within the configured session timeout, it assumes that the consumer has died. This can happen if the consumer’s JVM exits or experiences a long stop-the-world garbage collection pause.
• The Kafka broker may believe that a consumer is stuck and trigger rebalancing. If the consumer takes longer than a threshold to poll for the next batch of records, the broker marks it as stuck. This can happen when processing the previous batch takes too long. ([View Highlight](https://read.readwise.io/read/01m3nkw141bz0wtqtxr26p3xx8))
- Consumer rebalancing is needed to ensure partitions are evenly distributed. However, rebalancing can cause disruption and increased latency, particularly due to the near real-time nature of the e-commerce landscape. ([View Highlight](https://read.readwise.io/read/01m3nkw63gp5qwmxb1j0zjqwb4))
- 2 - Poison Pill Messages
A “poison pill” message in Kafka is a message that consistently causes a consumer to fail when attempting to process it ([View Highlight](https://read.readwise.io/read/01m3nkwpkq0jrf3sx3ybj7ekw7))
- When the consumer encounters such a message, it will fail to process it and throw an exception. By default, the consumer will return to the broker to fetch the same batch of messages again. Since the poison pill message is still present in that batch, the consumer will again fail to process it, and this loop continues indefinitely.
As a result, the consumer gets stuck on this one bad message and is unable to make progress on other messages in the partition. This is similar to the “head-of-line blocking” problem in networking. ([View Highlight](https://read.readwise.io/read/01m3nkx9z2epwvf2k9cpz2x90f))
- 3 - Cost Concerns
There is a strong coupling between the number of partitions in a Kafka topic and the maximum number of consumers that can read from that topic in parallel. This coupling can lead to increased costs when trying to scale consumer applications to handle higher throughput.
For example, consider that you have a Kafka topic with 10 partitions and 10 consumer instances reading from this topic. Now, if the rate of incoming messages increases and the consumers are unable to keep up (i.e. consumer lag starts to increase), you might want to scale up your consumer application by adding more instances.
However, once you have 10 consumers (one for each partition) in a single group, adding more consumers to that group won’t help because Kafka will not assign more than one consumer from the same group to a partition. The only way to allow more consumers in a group is to increase the number of partitions in the topic. ([View Highlight](https://read.readwise.io/read/01m3nkzv3nk0ve3htc90t7ncwh))
- Walmart has a massive Apache Kafka deployment with 25K+ consumers across private and public cloud environments.
This deployment processes trillions of Kafka messages per day at 99.99% availability. It supports critical use cases such as:
• Movement of data
• Event-driven microservices
• Streaming analytics ([View Highlight](https://read.readwise.io/read/01m3nkf7qf08szzat42npks30z))
- However, increasing the number of partitions comes with its challenges and costs.
• Kafka has a recommended limit on the number of partitions per broker (for example, 4000 partitions per broker). If you keep increasing partitions, you may hit this limit and need to scale the Kafka brokers to larger instances, even if the brokers have sufficient resources to handle the current load. Scaling to larger broker instances is expensive.
• Increasing partitions requires coordination among the Kafka team, the producer, and the consumer teams. In a large organization with thousands of Kafka pipelines, this coordination overhead is significant.
• More partitions also mean more open file handles, increased memory usage, and more threads on the Kafka brokers. This can lead to higher resource utilization and costs. ([View Highlight](https://read.readwise.io/read/01m3nm0fmn8c9gtfrfv99ypxv7))
- At Walmart's scale, the Kafka setup must be able to handle sudden traffic spikes. Also, consumer applications are written in multiple languages. Therefore, all consumer applications must adopt some best practices to maintain the same level of reliability and quality. ([View Highlight](https://read.readwise.io/read/01m3nkffznzg0gvqh6n58w8gzg))
- Challenges with Kafka at Walmart’s Scale
Let’s start with understanding the main challenges that Walmart faced. ([View Highlight](https://read.readwise.io/read/01m3nkfwr2kh65kw4tk5tdb24e))
- One of the most frequent problems was related to consumer rebalancing.
But what triggers consumer rebalancing in Kafka?
This can happen due to the changing number of consumer instances within a consumer group. ([View Highlight](https://read.readwise.io/read/01m3nkhcc4k3s8c0jjjqye52j8))