How do you implement a custom Partitioner in Kafka, and what is a legitimate use case for one?
answer
- implement Partitioner.partition(...)
- partitioner.class = MyPartitioner
- Cluster.availablePartitionsForTopic
- hot path: fast, allocation-light, thread-safe
- use case: hot-key isolation / tenant routing
basics
~10 sImplement 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 sYou 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
Know a custom partitioner is possible via the Partitioner interface and partitioner.class, even if not the details.
Name the partition() signature and registration config and give one real use case.
Discuss the Cluster API, hot-path/thread-safety/determinism constraints, and skew-control routing.
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.