skip to content

How do you implement a custom Partitioner in Kafka, and what is a legitimate use case for one?

level: seniorimportance: should knowfreq 42%

answer

  1. implement Partitioner.partition(...)
  2. partitioner.class = MyPartitioner
  3. Cluster.availablePartitionsForTopic
  4. hot path: fast, allocation-light, thread-safe
  5. use case: hot-key isolation / tenant routing

basics

~10 s

Implement org.apache.kafka.clients.producer.Partitioner (the partition() method returns an int partition), then set partitioner.class to your class. A common reason is routing hot or special keys to dedicated partitions to control skew or isolate VIP traffic.

solid answer

~50 s

You implement `org.apache.kafka.clients.producer.Partitioner`, whose key method is `int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster)`. You use the `Cluster` to read partition counts and which partitions have a live leader (`cluster.availablePartitionsForTopic`), then return the chosen partition index. You also implement `configure(Map)` and `close()` (it extends `Configurable`/`Closeable`). Register it via `props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, MyPartitioner.class)`. Legitimate use cases: pinning known hot keys to a dedicated partition to bound skew; geo/tenant routing; sending a priority tenant to reserved partitions; or replicating an external system's partitioning scheme for co-partitioning. Pitfalls: `partition()` is on the hot send path so it must be fast and allocation-light; it must be thread-safe (one producer, many app threads); and any custom hashing must be deterministic or you lose per-key ordering. Note KIP-794 made the built-in load-aware partitioner the default direction, so reach for custom only when routing logic is genuinely domain-specific.

go deeper

for a junior

Know a custom partitioner is possible via the Partitioner interface and partitioner.class, even if not the details.

for a middle

Name the partition() signature and registration config and give one real use case.

for a senior

Discuss the Cluster API, hot-path/thread-safety/determinism constraints, and skew-control routing.

for a principal

Weigh custom routing against the built-in load-aware partitioner, design tenant/priority partitioning schemes, and own the operational risks.

## The Partitioner interface Kafka lets you override partition selection by implementing `org.apache.kafka.clients.producer.Partitioner`. It extends `Configurable` and `Closeable`, so you implement three methods: - `void configure(Map<String,?> configs)` — called once at producer startup with the producer configs; read your custom settings here. - `int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster)` — returns the target partition index. Called for **every record**. - `void close()` — cleanup. You register it with `partitioner.class`: `props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, MyPartitioner.class.getName());` ## Using the Cluster argument The `Cluster` object exposes topic metadata: `cluster.partitionCountForTopic(topic)` and `cluster.availablePartitionsForTopic(topic)` (partitions that currently have a leader). Routing to only *available* partitions avoids blocking on a partition whose leader is down. `cluster.partitionsForTopic(topic)` returns all partitions including offline ones. ## Legitimate use cases 1. **Hot-key isolation / skew control:** if a handful of keys (a giant tenant) would overload their hashed partition, route those specific keys to dedicated partitions and hash the rest normally, bounding skew. 2. **Tenant / geography routing:** map a tenant id or region to a partition range so each consumer instance owns a region. 3. **Priority lanes:** reserve partitions for high-priority traffic so a separate consumer group services them with lower latency. 4. **Co-partitioning with an external system:** mirror another system's hashing so joins line up. ## Pitfalls and constraints - **Hot path:** `partition()` runs synchronously on every `send()`. Keep it O(1), avoid per-call allocations, locks, or I/O — never call a DB or remote service inside it. - **Thread safety:** a single `KafkaProducer` is shared by many application threads, all calling `partition()` concurrently. Any mutable state (e.g. a round-robin counter) must be thread-safe (e.g. `AtomicInteger`). - **Determinism for ordering:** if you intend per-key ordering, the same key must always map to the same partition; non-deterministic logic (random, time-based) breaks ordering. - **Partition-count changes:** like the default, your modulo-style logic is sensitive to partition count changes; design for it. - **Validity:** you must return a partition that exists and ideally has a leader; returning an offline partition stalls sends. ## Modern context KIP-794 (Kafka 3.3) improved the *built-in* partitioner to be load-aware and deprecated `DefaultPartitioner`/`UniformStickyPartitioner`. So a custom partitioner is justified mainly when your routing is **domain-specific** (tenant/geo/priority/hot-key), not merely to chase uniform load — the platform now does uniform-load balancing for you.

  • Why must a custom partitioner avoid blocking work inside partition()?
    partition() is invoked synchronously on every send() on the producer's hot path. Any lock contention, allocation, or I/O (e.g. a DB lookup) directly adds latency to every produce call and can throttle throughput across all threads sharing the producer.
  • What thread-safety concern arises if your custom partitioner keeps a round-robin counter?
    One KafkaProducer is shared by many threads that all call partition() concurrently, so a plain int counter would race. Use an AtomicInteger (or other concurrent construct) to increment safely.

saying these in an interview costs you the question

  • Doing a remote/DB lookup inside partition() (blocks the hot send path).
  • Keeping non-thread-safe mutable state in the partitioner.
  • Using random/time-based logic while claiming per-key ordering.
  • Returning a partition index that may not exist or has no leader.
  • Writing a custom partitioner just to balance load now that KIP-794's default does it.

context