Kafka in Production: Partition Strategy, Consumer Lag, Reliability, and an Incident Playbook

A production Kafka cluster needs more than brokers. Learn to choose partition keys, plan retention and capacity, monitor lag, handle rebalances, and rehearse failure recovery.
A Kafka demo succeeds when a producer sends a record and a consumer prints it. A production system succeeds when a broker fails, a consumer is slow, a schema changes, and the business result remains understandable. The hard questions are operational: what must stay ordered, how much backlog can accumulate, who owns an event contract, and how will you recover after a bad release?
This guide turns those questions into a practical design and review process. The example is an order platform with orders.placed.v1, an inventory consumer group, and an analytics group. The numbers are illustrative; measure your own payloads, throughput, and recovery target before choosing a cluster size.
Start with the workload, not a broker count
Write down peak events per second, average and high-percentile record size, retention period, required read fan-out, acceptable processing delay, and recovery time after one consumer or broker fails. Identify whether every order's events must stay in order and whether all orders must share one total order. Most systems need per-entity ordering, not a global order that would force one partition.
For example, 2,000 events per second at an average 1 KB is roughly 2 MB/s of raw payload before replication, protocol overhead, indexes, consumers, and bursts. Seven days of raw retained payload is already over a terabyte by rough decimal arithmetic. This is a planning estimate, not a disk-sizing formula. Compression, retention settings, replicas, actual record sizes, and peak traffic change the result. Benchmark and leave headroom.
Avoid planning only for average throughput. A promotion or bulk import can create a hot period, and a consumer outage creates a backlog that must be drained faster than new events arrive. If the consumer normally handles exactly the incoming rate, it never catches up after an outage.
Choose partition keys for the ordering you need
Kafka orders records within each topic partition. Use orderId as the key if events for the same order must be processed in order. customerId changes the promise to per-customer ordering. A constant key concentrates all traffic on one partition; a random key distributes load but loses related-event order.
Partition count limits parallel active consumers within one group for that topic. Eight partitions can be assigned to at most eight active consumers in that group for that topic. It is not a promise that eight consumers will have equal work: a hot key can overload one partition while others are quiet.
Increasing partitions later can alter which partition receives a key for future events. If strict long-lived ordering by key matters, plan the migration deliberately, perhaps with a new topic and a controlled cutover. Measure skew per partition rather than trusting topic-level averages. Do not create hundreds of partitions as a reflex; they consume cluster resources and make operations more complex.
Set replication and acknowledgement expectations
A replication factor specifies how many copies of a partition the cluster maintains. A production design often uses three replicas across suitable failure domains, but availability depends on placement, in-sync replicas, broker health, producer acknowledgements, and min.insync.replicas. Review these settings together. acks=all without a meaningful minimum in-sync replica policy does not express the entire durability requirement.
An idempotent producer reduces duplicates caused by producer retries to Kafka. It does not deduplicate two independent business requests or make an external database transaction atomic with the publish. Kafka transactions can coordinate Kafka output records and consumed offsets for Kafka-to-Kafka processing; side effects in PostgreSQL, payment providers, or email need their own strategy.
Test the precise failure you claim to tolerate. Shut down a broker in a staging cluster with representative replication settings; verify producer errors, recovery time, consumer behavior, and whether your service degrades gracefully. A single-broker local Docker setup is excellent for learning, but it proves nothing about broker failover.
Kafka 4.x operates in KRaft mode; ZooKeeper mode was removed in Kafka 4.0. Treat older guides that require a ZooKeeper ensemble as historical material. Keep controller quorum design, backup, upgrades, and broker operations aligned with the release you actually deploy.
Treat consumer progress as a business signal
A consumer group commits offsets to record its position. Lag compares the latest available offset with the group's position per partition. It is an important clue, but offset lag is not identical to wall-clock delay. A thousand tiny records may clear quickly; one slow call per record may take hours. Track both offset lag and the age of the oldest unprocessed business event.
Useful signals include:
| Signal | What it suggests | First question |
|---|---|---|
| Lag grows across all partitions | Consumer throughput below input or dependency slow | Did input spike or processing slow? |
| One partition lags | Hot key, skew, stuck record | Which key or offset is blocking it? |
| Frequent rebalances | Membership churn or processing/poll issues | Are pods restarting or polls delayed? |
| Outbox age grows | Producer pipeline stuck before Kafka | Are publishers failing or unclaimed? |
| Under-replicated partitions | Replica cannot keep up or broker unavailable | Is a broker, disk, or network unhealthy? |
| Dead-letter volume rises | Poison records or incompatible schema | Which event version first failed? |
Alert on sustained violations tied to business tolerance. If fulfillment must react within two minutes, measure event-to-effect time rather than only a generic broker metric. Include partition and group labels in dashboards, but avoid exploding metric cardinality with order IDs.
Understand rebalances and slow processing
When group membership or topic metadata changes, partitions can be reassigned. A consumer that holds an in-memory batch should not assume continued ownership after revocation. Use your client's documented rebalance hooks when needed, coordinate in-flight work, and commit only completed processing. Long blocking work in the poll loop can destabilize group membership; move it to a bounded worker pipeline with careful offset tracking if your client model supports that.
Scaling consumers helps only when there are partitions to assign and the bottleneck is consumer CPU or independent work. If each event waits on the same overloaded PostgreSQL table or external API, adding more consumers can worsen the downstream failure. Limit concurrency and use backpressure; consider partition pause/resume behavior where supported.
A consumer crash after a database write but before offset commit causes replay. Design idempotent writes keyed by eventId or another stable operation ID. If you commit before the side effect, a crash can lose work. There is no universal “exactly once” switch for arbitrary external effects.
Make retention and replay policy explicit
Retention defines how long Kafka keeps records, whether or not consumers have read them. It must exceed expected outages and the replay window you promise, with margin. If a group falls behind beyond the earliest retained offset, it cannot simply continue from data that is gone. Plan alerts before that point and a recovery source if history matters longer.
Compaction preserves a latest value per key over time for changelog-style topics; it is not a complete event audit trail. Use time-retained event topics for history, compacted topics for state where appropriate, and separate access policies. Deletion requirements for personal data need specific design because events can exist in replicas, backups, and downstream stores.
Replaying a projection can be safe if it rebuilds an isolated table. Replaying an email or charge is a different risk. Keep side-effect deduplication and a documented replay runbook: which offsets or timestamps, which group, what downstream effects are disabled, how progress is validated, and when normal processing resumes.
Govern event schemas and ownership
An event contract should identify its owner, meaning, key, schema version, required and optional fields, timestamps, and compatibility policy. Test producer changes against consumers that may be behind. An additive field can still break a strict decoder; coordinate client expectations. A changed meaning of status is more dangerous than a new field with a clear name.
Use a schema format and registry if they fit your ecosystem; JSON with validated schemas can also work for a smaller system. What matters is an enforceable compatibility rule and examples that CI checks. Include a stable eventId and distinguish occurrence time from processing time. Do not publish whole ORM entities or secrets.
Keep topic ownership clear. A team should know who approves retention changes, partition increases, schema changes, and access grants. Without ownership, a shared topic becomes a hidden coupling point for many services.
Secure the cluster and data
Use encryption in transit and appropriate authentication for clients, then grant topic and group permissions with least privilege. Rotate credentials and keep them in a secret store, not application config committed to Git. Separate development and production clusters or access scopes. Audit administrative operations and restrict who can alter retention or delete topics.
Treat event payloads as data leaving the source service. Minimize personal information, set access and retention accordingly, and avoid logging complete payloads. A producer with write access to one topic should not automatically read every consumer's topic. Security review belongs in topic onboarding, not after a breach.
A release and incident playbook
Before a release, validate schema compatibility, run a representative load test, and confirm dashboards for producer failures, outbox age, consumer lag, rebalance rate, and business completion delay. Rehearse a rollback that does not change the event meaning midway. If a new consumer is deployed, start it from a consciously chosen offset; do not leave earliest versus latest to an undocumented default.
When lag suddenly rises:
- Scope it. Which group, topic, and partitions? Did input rate rise or processing rate fall?
- Check recent changes. Deployments, schema changes, broker maintenance, dependency outages, and credentials.
- Inspect a safe sample. Find the failing offset and error category without exposing sensitive payloads.
- Protect downstream systems. Reduce concurrency or pause a partition if retries are amplifying an outage.
- Recover deliberately. Fix code or data, then resume or replay with idempotency checks.
- Verify business completion. Falling lag alone does not prove that reservations or notifications were correct.
If only one partition is stuck on a poison record, blindly adding consumers will not help. If every partition lags because the database is slow, extra consumers can compound contention. Separate the symptom from its bottleneck.
Production readiness checklist
- The partition key matches the required ordering scope and has measured distribution.
- Replication, producer acknowledgements, and failure domains match the durability target.
- Consumer side effects are safe on replay, and offsets advance after completed work.
- Retention exceeds expected outage and replay needs; storage has measured headroom.
- Event schema, ownership, compatibility, and access are documented.
- Dashboards show lag and event-to-business-effect delay.
- The team has rehearsed a broker failure, a consumer crash, a poison record, and a replay.
Kafka gives you a powerful distributed log. Production reliability comes from the decisions around it: event identity, partitioning, side effects, capacity, access, and recovery. Design those before the first traffic spike makes them urgent.
Official and vendor documentation
Featured Articles

Laravel and Kafka Without Lost Events: The Transactional Outbox, Idempotent Consumers, and PostgreSQL
A database commit and a Kafka publish cannot safely be treated as one ordinary transaction. This guide builds an outbox and consumer design that survives crashes and duplicate delivery.

Apache Kafka Explained: Topics, Partitions, Consumer Groups, and Your First Event Pipeline
Follow one order event from producer to consumers, then run a local Kafka topic and learn what partitions, offsets, keys, and consumer groups actually do.

How to Design APIs That Clients Can Trust: A Practical Contract-First Guide
Good APIs make the next client request predictable. Design an enrollment API from the use case outward, with clear contracts, safe retries, useful errors, and a plan for change.
Comments
0 commentsNo approved comments are visible yet. New community replies may wait for moderation.