skip to content

How do you configure partitioned producers and consumers in Spring Cloud Stream, and why would you partition?

level: seniorimportance: should knowfreq 45%

answer

  1. producer: partitionKeyExpression / partitionKeyExtractorName + partitionCount
  2. selector default = key.hashCode() % partitionCount
  3. consumer: partitioned=true + instanceCount/instanceIndex
  4. Kafka auto-assigns via rebalance; Rabbit needs instanceIndex
  5. same key → same partition → per-key order; watch hot partitions

basics

~20 s

Partitioning routes messages that share a key to the same partition, so one consumer instance always handles them in order. On the producer you set a key expression and partition count; on the consumer you enable partitioned consumption.

solid answer

~40 s

Partitioning guarantees that all messages with the same key land on the same partition, giving **per-key ordering** and consistent stateful processing. On the **producer** side you set `producer.partitionKeyExpression` (a SpEL expression over the outbound message) or a `partitionKeyExtractorName` bean, plus `producer.partitionCount`. A `partitionSelectorExpression`/`partitionSelectorName` maps the key to a partition index; the default is `key.hashCode() % partitionCount`. On the **consumer** side you set `consumer.partitioned=true` and, for binders without native partition management, the top-level `spring.cloud.stream.instanceCount` and `instanceIndex` so each instance owns a slice of partitions. With the **Kafka** binder, native consumer-group rebalancing assigns partitions automatically, so `instanceIndex` matters less, but the producer key still drives which partition — and thus which instance — a message hits. Partition count caps consumer parallelism. Use partitioning when you need ordering per entity (per user, per account) or partition-local state.

code

java · 32 lines
java
// Producer + consumer partitioning (application.yml reference):
//
// # Producer: partition by accountId across 4 partitions
// spring.cloud.stream.bindings.emit-out-0.producer:
//   partition-key-expression: payload.accountId
//   partition-count: 4
//
// # Consumer: participate in partitioned consumption
// spring.cloud.stream.bindings.handle-in-0.consumer.partitioned: true
// spring.cloud.stream.instance-count: 4
// spring.cloud.stream.instance-index: 0   # unique per instance (0..3), e.g. via env

import java.util.function.Consumer;
import org.springframework.context.annotation.Bean;
import org.springframework.integration.partition.PartitionKeyExtractorStrategy;
import org.springframework.messaging.Message;

public class PartitioningExample {

    // Alternative to partition-key-expression: a strategy bean referenced by
    // producer.partition-key-extractor-name: accountKeyExtractor
    @Bean
    public PartitionKeyExtractorStrategy accountKeyExtractor() {
        return (Message<?> message) -> ((AccountEvent) message.getPayload()).accountId();
    }

    // All events for a given accountId are handled in order by one instance.
    @Bean
    public Consumer<AccountEvent> handle(Ledger ledger) {
        return ledger::apply;
    }
}

go deeper

for a junior

Grasp the idea: same key → same partition → one consumer handles it in order.

for a middle

List the producer key/count props and the consumer partitioned flag.

for a senior

Explain the selector default, instanceCount/instanceIndex vs Kafka native rebalance, and the parallelism ceiling.

for a principal

Reason about skew mitigation, custom selector strategies, repartitioning migration, and partition-blocking interaction with retry/DLQ.

**Why partition.** A consumer *group* spreads load, but by default which instance gets a given message is arbitrary — bad when messages for the same entity must be processed **in order** (e.g. all events for account #42) or when an instance keeps **local state** for a key. **Partitioning** solves this: messages are deterministically bucketed by a **partition key** so that all messages with the same key go to the same **partition**, and each partition is consumed by a single instance in the group. Result: **per-key ordering** and sticky routing. **Producer configuration** (`spring.cloud.stream.bindings.<out>.producer.*`): - **`partitionKeyExpression`** — a SpEL expression evaluated against the outbound `Message` to extract the key, e.g. `payload.accountId` or `headers['customerId']`. - **`partitionKeyExtractorName`** — alternatively, the bean name of a `PartitionKeyExtractorStrategy` implementation (use when key logic is too complex for SpEL). - **`partitionCount`** — how many partitions to spread across (the binder provisions this many for the destination, subject to broker limits). - **`partitionSelectorExpression`** / **`partitionSelectorName`** — how the extracted key maps to a partition *index*. Default algorithm: `key.hashCode() % partitionCount`. Provide a `PartitionSelectorStrategy` bean for custom mapping. Exactly one of `partitionKeyExpression` or `partitionKeyExtractorName` enables producer-side partitioning. **Consumer configuration:** - **`spring.cloud.stream.bindings.<in>.consumer.partitioned=true`** — declares this consumer participates in partitioned consumption. - **`spring.cloud.stream.instanceCount`** (top-level) — total number of consuming instances. - **`spring.cloud.stream.instanceIndex`** (top-level, per instance) — this instance's 0-based index; the binder uses it to decide which partitions this instance owns. These `instanceCount`/`instanceIndex` settings are essential for binders like **RabbitMQ** that don't natively manage partition assignment — Spring computes the assignment from them. For **Kafka**, native **consumer-group rebalancing** already assigns partitions dynamically, so you often don't need to hand-set `instanceIndex`; but you still set `partitionCount` on the producer (or provision topic partitions) and a partition key so routing is deterministic. Kafka's own message key (set from the partition key) drives native partitioning. **Parallelism ceiling.** With N partitions, at most N instances in the group actively consume; extras idle. Choose `partitionCount` for your target scale plus headroom. **Edge cases & gotchas.** - **Hot partitions / skew:** if one key (or a few) dominates traffic, its partition is a bottleneck while others sit idle. Choose a high-cardinality, evenly distributed key. - **`hashCode()` stability:** the default selector uses Java `hashCode`. Non-deterministic or poorly distributed hashCodes cause uneven or shifting assignment — prefer stable keys (strings, ids) and a custom `PartitionSelectorStrategy` if needed. - **Repartitioning:** changing `partitionCount` later reshuffles key→partition mapping, breaking ordering continuity and possibly stateful assumptions. - **Ordering is per partition only** — there is no global order across partitions. - **Null key:** if the key expression yields null, routing falls back to round-robin/default and you lose the ordering guarantee. **Interaction with error handling.** Because retry is blocking per partition, a poison message stalls its *whole partition* (all keys sharing it), not just one key — a reason to pair partitioning with a DLQ so bad messages exit quickly. **When to use.** Partition when you need per-entity ordering or partition-local state/caching; skip it when messages are independent and you only want throughput (a plain group suffices).

  • Your traffic is dominated by one very active account, so one instance is overloaded while others idle. What's happening and how do you address it?
    Partition skew / a hot partition: that account's key always hashes to one partition. Options: pick a finer-grained key (e.g. account+region) if ordering allows, use a custom PartitionSelectorStrategy to spread hot keys, or relax the per-key ordering requirement so plain load-balancing applies.
  • Does partitioning give you global ordering of all messages?
    No — ordering is guaranteed only within a single partition (per key). Across partitions there is no global order, since they're consumed concurrently by different instances.

saying these in an interview costs you the question

  • Claiming partitioning provides a global total order across all messages
  • Thinking you can scale consumers beyond partitionCount and still get more parallelism
  • Ignoring key skew, producing hot partitions
  • Believing changing partitionCount later preserves the existing key→partition mapping

context