Why this matters
- VaultCommerce rolling deploys triggered 30-second rebalance storms with eager protocol — fixed with cooperative sticky.
- During rebalance, partitions may be revoked before reassignment — handle revocation gracefully.
- Rebalance listener hooks let VaultCommerce flush in-flight work before partition loss.
- Frequent rebalances often indicate slow poll loops or unstable membership.
Eager vs cooperative
Eager (range, round-robin): all partitions revoked, then reassigned — full stop. Cooperative sticky: only moved partitions revoked incrementally. VaultCommerce sets partition.assignment.strategy=CooperativeStickyAssignor on all Java consumers.
Key points
- Rebalance — partition ownership transfer among group members
- CooperativeStickyAssignor — incremental rebalance (Kafka 2.4+)
- onPartitionsRevoked — callback to commit offsets before losing partitions
- group.instance.id — static member identity across restarts
- max.poll.interval.ms — max time between poll() calls before considered dead
Triggers and mitigation
Join, leave, session timeout, max.poll.interval exceeded. Use group.instance.id for static membership so restarts don't look like new members. Increase session.timeout.ms only with proportional heartbeat.interval.ms.
VaultCommerce rollout checklist
Before promoting changes that touch inventory-service and payment-service consumer groups, run the staging KRaft cluster (Kafka 3.7+, Schema Registry 7.x) through a 10k events/min soak test. Compare producer request latency p99 and consumer lag per group against the pre-deploy baseline. Use CooperativeStickyAssignor in production. Document the change in the internal topic registry, attach Grafana screenshots to the change ticket, and keep an engineer on lag dashboards for 30 minutes after production rollout — roll back the service release before altering broker-level settings if lag or under-replicated partitions spike.
// VaultCommerce — cooperative rebalance + revocation handler
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
CooperativeStickyAssignor.class.getName());
props.put(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, "inventory-pod-" + podName);
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 45_000);
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 15_000);
consumer.subscribe(List.of("vaultcommerce.orders.placed.v1"), new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
consumer.commitSync(); // flush offsets before handoff
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> partitions) { }
});
Quick recall
Everything you need if you only revisit this box.
- Use CooperativeStickyAssignor in production.
- Commit on partition revocation.
- Static group.instance.id reduces deploy churn.
Test yourself
Answer these before moving on — recall is what makes it stick.
- How do you perform rebalancing and what are the pitfalls (stability, consumer downtime)?
- Describe a real-world example: designing an event-driven order processing pipeline with Kafka—what are the key design decisions?
- Explain the core concepts of Kafka Streams: KStream, KTable, state stores, and joins.