skip to content

What is the ReplicaSelector interface, and how could you implement a custom replica selector beyond RackAwareReplicaSelector?

level: seniorimportance: nice to knowfreq 25%

answer

  1. org.apache.kafka.common.replica.ReplicaSelector
  2. select(tp, ClientMetadata, PartitionView) -> ReplicaView
  3. built-ins: LeaderSelector, RackAwareReplicaSelector
  4. runs on hot fetch path — keep it fast
  5. deploy JAR on broker classpath

basics

~20 s

replica.selector.class plugs in a class implementing org.apache.kafka.common.replica.ReplicaSelector. Kafka ships LeaderSelector (default) and RackAwareReplicaSelector. You can write your own select() logic — e.g., choose by latency or a custom topology — returning the replica the consumer should read from.

solid answer

~50 s

ReplicaSelector is the broker-side extension point KIP-392 added. The broker passes the selector the set of replica views for a partition and the client's metadata (including client.rack), and select() returns the ReplicaView the consumer should be steered to (surfaced as preferredReadReplica). The built-ins are LeaderSelector (always the leader) and RackAwareReplicaSelector (in-rack ISR replica, else leader). To customize, implement ReplicaSelector: in select(), inspect the ClientMetadata (rack, client id, address) and each ReplicaView (replica info plus its log end offset / last-caught-up time), apply your policy — for instance prefer the lowest-latency or least-loaded in-sync replica, or honor a richer topology than a flat rack string — and return the chosen replica, defaulting to the leader when no candidate fits. You also implement configure() and close(). Package it on the broker classpath and point replica.selector.class at it.

code

java · 15 lines
java
public class LatencyAwareReplicaSelector implements ReplicaSelector {
    @Override
    public Optional<ReplicaView> select(TopicPartition tp,
                                        ClientMetadata metadata,
                                        PartitionView partitionView) {
        String clientRack = metadata.rackId();
        return partitionView.replicas().stream()
            // only in-sync, caught-up replicas are safe candidates
            .filter(r -> isEligible(r, partitionView))
            .min(Comparator
                .comparingInt((ReplicaView r) -> sameRack(r, clientRack) ? 0 : 1)
                .thenComparingLong(this::estimatedLatencyMs))
            .or(() -> Optional.of(partitionView.leader()));
    }
}

go deeper

for a junior

Know there are two built-in selectors and that the rack-aware one enables FFF; custom ones are advanced.

for a middle

Identify ReplicaSelector as a broker plugin and name the two built-ins.

for a senior

Describe the select() signature, the ClientMetadata/ReplicaView inputs, and a realistic custom policy (latency/topology) with hot-path cautions.

for a principal

Judge when a custom selector is warranted vs over-engineering, and design one against cost/latency/freshness objectives safely.

## The extension point KIP-392 made replica selection **pluggable** via the broker config `replica.selector.class`, which names an implementation of `org.apache.kafka.common.replica.ReplicaSelector`. The broker, when answering a consumer fetch, consults this selector to decide which replica the consumer *should* read from and returns that as the **preferred read replica**. ## The interface, conceptually `ReplicaSelector` provides (roughly): - `select(TopicPartition tp, ClientMetadata metadata, PartitionView partitionView)` returning an `Optional<ReplicaView>` — the chosen replica, or empty/leader fallback. - `configure(Map<String,?> configs)` — pick up custom config. - `close()` — cleanup. Supporting types: - **ClientMetadata** — the requesting consumer's info: `rackId` (its `client.rack`), client id, listener, and address. Your policy keys off this. - **ReplicaView / PartitionView** — each candidate replica's info: which broker, its **log end offset**, and **time of last caught-up** (so you can gauge how fresh/in-sync it is); the partition view exposes the leader and the set of replicas. ## Built-in implementations - **`LeaderSelector`** (default): always returns the leader → classic, FFF off. - **`RackAwareReplicaSelector`**: among ISR replicas, prefer one whose `broker.rack` equals the client's `rackId`; else the leader. ## Writing a custom selector Reasons you might: a flat rack string can't express your real topology (e.g., region + AZ + cost tiers), or you want **latency-based** or **load-based** steering, or you want to bias toward replicas with the freshest HW. Sketch: 1. `configure()` — read any custom properties (e.g., a latency map). 2. `select()` — filter to eligible (in-sync, has the needed offset) replicas using the ReplicaViews; rank them by your metric (rack match, then latency, then freshness); return the best, or `Optional.empty()`/leader if none qualifies. 3. Keep it **fast and side-effect-free** — it runs on the hot fetch path for every relevant fetch. 4. Deploy the JAR on every broker's classpath; set `replica.selector.class` to your class; roll it out. ## Practical cautions - The selector only *suggests*; the consumer still validates offsets and may fall back if the chosen replica can't serve the requested offset. - A buggy or slow selector affects fetch latency cluster-wide — test thoroughly. - Most teams never need a custom one; RackAwareReplicaSelector covers the cost-driven AZ-locality case. Custom selectors are for unusual topologies or cost models.

  • What information does the selector get about the requesting consumer?
    ClientMetadata: its rackId (client.rack), client id, the listener it connected on, and its address — enough to apply rack- or location-based policies.
  • Why must a custom selector be fast and side-effect-free?
    select() runs on the broker's fetch hot path for affected fetches; slow or blocking logic would add latency to consumer reads cluster-wide.

saying these in an interview costs you the question

  • Saying the selector runs on the consumer — it's a broker-side plugin.
  • Thinking RackAwareReplicaSelector is the only possible implementation — the interface is pluggable.
  • Assuming the selector's choice is binding regardless of offset availability — the consumer still validates and can fall back.

context