skip to content

What is the role of KafkaApis in the request pipeline, and how does it decide what to do with a request?

level: middleimportance: should knowfreq 35%

answer

  1. KafkaApis.handle() = dispatch on ApiKeys
  2. handleProduceRequest / handleFetchRequest / handleMetadataRequest ...
  3. authz + validate + delegate
  4. ReplicaManager, GroupCoordinator, metadata cache, TxnCoordinator
  5. callbacks/purgatory build the response later

basics

~10 s

KafkaApis is the broker class that actually handles requests. A request-handler thread calls KafkaApis.handle(), which looks at the request's API key (Produce, Fetch, Metadata, etc.) and dispatches to the matching handler method.

solid answer

~40 s

KafkaApis is the request-dispatch and business-logic entry point on a broker. After a request-handler (I/O) thread dequeues a request from the request queue, it calls KafkaApis.handle(request). Internally that is a switch on the request's ApiKeys value (PRODUCE, FETCH, METADATA, OFFSET_COMMIT, LIST_OFFSETS, JOIN_GROUP, …) routing to a specific handleXxxRequest method. Each method does authorization checks, validates the request, and delegates to the right subsystem: ReplicaManager for produce/fetch, GroupCoordinator for consumer-group APIs, the metadata cache for metadata, etc. It also builds and sends the response (often via a callback once asynchronous work or purgatory completes). So KafkaApis is the 'controller layer' that maps the wire protocol's API keys onto the broker's internal components.

go deeper

for a junior

KafkaApis is the class that handles requests, choosing logic based on the request type.

for a middle

Explain dispatch by ApiKeys to handleXxxRequest and delegation to ReplicaManager/GroupCoordinator/metadata cache.

for a senior

Add where authorization/validation live and how purgatory/callbacks make handling asynchronous.

for a principal

Discuss KafkaApis as the protocol/internals seam and how KIPs extend it with new ApiKeys and handlers.

`KafkaApis` is the broker component that turns a parsed protocol request into actual work — think of it as the broker's request router and controller layer. **Entry point.** Every request-handler / I/O thread, after pulling a request off the shared request queue, calls **`KafkaApis.handle(request)`**. There is effectively one logical `KafkaApis` instance per broker shared across handler threads. **Dispatch by API key.** The Kafka wire protocol assigns each request type a numeric **API key**, enumerated in **`ApiKeys`** (`PRODUCE=0`, `FETCH=1`, `LIST_OFFSETS=2`, `METADATA=3`, `OFFSET_COMMIT=8`, `JOIN_GROUP=11`, `API_VERSIONS=18`, …). `handle()` reads `request.header.apiKey` and switches to the matching **`handleXxxRequest`** method — `handleProduceRequest`, `handleFetchRequest`, `handleMetadataRequest`, `handleJoinGroupRequest`, and so on. **What each handler does.** 1. **Authorization** — checks ACLs via the `Authorizer` for the resource/operation. 2. **Validation** — version compatibility, topic/partition existence, request well-formedness. 3. **Delegation** to the right subsystem: - Produce/Fetch ⇒ **`ReplicaManager`** (append to / read from logs, manage replication). - Consumer-group APIs (JoinGroup, SyncGroup, OffsetCommit, OffsetFetch, Heartbeat) ⇒ **`GroupCoordinator`**. - Metadata ⇒ the broker's **metadata cache**. - Transactions ⇒ **`TransactionCoordinator`**. 4. **Response construction** — builds the response object and sends it back through the `RequestChannel`, which routes it to the originating Processor's response queue. **Asynchrony and purgatory.** Many handlers don't complete inline. A fetch that needs to wait for `fetch.min.bytes`, or a produce with `acks=all` waiting on replicas, registers a delayed operation in a **purgatory** and returns; a callback fires later to build and enqueue the response. This keeps handler threads free. **Why it matters.** `KafkaApis` is where protocol concerns (API keys, versions) meet broker internals (replica manager, coordinators). Understanding it explains how a single class can implement the entire broker-side surface of the Kafka protocol and where authorization and version negotiation live. It is also the seam most relevant when reasoning about new protocol APIs introduced by KIPs — each new request type adds an `ApiKeys` entry and a `handleXxxRequest` branch.

  • Which subsystem does KafkaApis delegate a FETCH request to?
    ReplicaManager (ReplicaManager.fetchMessages), which reads from the partition logs / page cache and may park the request in the fetch purgatory until fetch.min.bytes or the timeout is satisfied.
  • Where in the request handling does ACL authorization happen?
    Inside the KafkaApis handler method for that request, via the configured Authorizer, before the request is allowed to act on the resource.

saying these in an interview costs you the question

  • Saying KafkaApis runs on the network/processor threads (it runs on the request-handler/I/O threads).
  • Thinking there's a separate class per request type instead of dispatch by ApiKeys inside KafkaApis.
  • Believing all handlers complete synchronously (many use callbacks/purgatory).

context