skip to content

Describe the three internal topics a distributed Connect cluster uses (config, offset, status). What does each store, and what configuration do they require?

level: seniorimportance: must knowfreq 60%

answer

  1. config = 1 partition, total order
  2. offset = source positions, ~25 parts
  3. status = RUNNING/FAILED, REST reads it
  4. all compacted
  5. sink offsets → __consumer_offsets

basics

~10 s

config.storage.topic holds connector/task configs (single partition, compacted). offset.storage.topic holds source-connector read positions (many partitions, compacted). status.storage.topic holds connector/task/worker status (many partitions, compacted). All need high replication.

solid answer

~40 s

Distributed Connect persists all cluster state in three compacted Kafka topics, named via worker config. `config.storage.topic` stores every connector and task configuration; it must have exactly **one partition** so config changes are totally ordered, and `config.storage.replication.factor` should be 3+ in prod. `offset.storage.topic` stores **source** connector offsets (the framework-managed position in the external system, e.g., a DB SCN or file byte offset); use multiple partitions (`offset.storage.partitions`, often 25) and high replication. `status.storage.topic` stores the live state of connectors/tasks (RUNNING/FAILED/PAUSED/UNASSIGNED) surfaced by the REST API; multiple partitions, high replication. All three are **log-compacted** so the latest value per key survives indefinitely. Sink offsets are NOT in offset.storage.topic — sinks commit to `__consumer_offsets` like any consumer group.

go deeper

for a junior

Just name the three topics and that they store config/offset/status in Kafka.

for a middle

Know what each stores and that they are compacted with high replication.

for a senior

Explain the single-partition config requirement and the source-vs-sink offset distinction.

for a principal

Own the topic provisioning policy: partitions/RF/compaction, multi-cluster isolation, and failure modes from misconfiguration.

## The three internal topics In **distributed mode**, a Connect cluster keeps no durable state on local disk — everything lives in Kafka so any worker can recover it after a rebalance. Three topics, named by worker properties, hold that state. ### 1. `config.storage.topic` - **Stores**: connector configurations and the derived task configurations. - **Partitions**: **must be exactly 1**. Config is a single totally-ordered log so every worker, replaying it, reaches the identical final state. Multiple partitions would break ordering guarantees. - **Replication**: `config.storage.replication.factor` ≥ 3 in production. - **Cleanup policy**: **compact** — keep the latest config per key forever. ### 2. `offset.storage.topic` - **Stores**: **source** connector offsets — an opaque, connector-defined position in the *source* system (a file byte offset, a database log sequence number, an API cursor). This is how a source resumes exactly where it left off after restart/rebalance. - **Partitions**: `offset.storage.partitions`, commonly **25**, to spread offset writes. - **Replication**: `offset.storage.replication.factor` ≥ 3. - **Cleanup policy**: **compact**. - **Important**: sink connectors do **not** use this topic. A sink is a consumer group, so it commits its progress to the cluster's `__consumer_offsets` topic under its consumer group id. ### 3. `status.storage.topic` - **Stores**: the operational status of connectors, tasks, and workers — states like `RUNNING`, `PAUSED`, `FAILED` (with the failure trace), `UNASSIGNED`. The REST API `GET /connectors/{name}/status` reads from here. - **Partitions**: `status.storage.partitions`, commonly **5**. - **Replication**: `status.storage.replication.factor` ≥ 3. - **Cleanup policy**: **compact**. ### Why compaction everywhere? Compaction retains the most recent record per key and garbage-collects older ones, so the topics act as durable key-value stores of 'current state' rather than ever-growing logs. If a topic were created with `cleanup.policy=delete` (e.g., wrong auto-create defaults), old configs/offsets could be deleted and the cluster would lose state — a classic production incident. ### Operational notes / edge cases - Connect can auto-create these topics, but in locked-down clusters you **pre-create** them with correct partitions/replication/compaction. - The **same internal topic names must NOT be shared** by two Connect clusters with different `group.id`s — their state would collide and corrupt each other. - Replication factor of 1 (common in single-broker dev) means losing that broker loses all cluster state — never do this in prod. - A frequent confusion: 'how many partitions for config?' Always **1**.

  • Why must config.storage.topic have exactly one partition?
    To guarantee a total order of configuration changes. Every worker replays the single log and converges to the same final config; multiple partitions would lose that global ordering.
  • Where do SINK connector offsets get committed?
    Sinks behave as a consumer group and commit to the cluster's __consumer_offsets topic under their consumer group id — not to offset.storage.topic, which is only for source offsets.
  • What happens if these internal topics are created with cleanup.policy=delete instead of compact?
    Old keys (configs/offsets/status) can be deleted by retention, so the cluster loses durable state — connectors may lose their configs or resume position. They must be compacted.

saying these in an interview costs you the question

  • Saying config.storage.topic should have many partitions for throughput (it must be exactly 1).
  • Claiming sink offsets live in offset.storage.topic (they go to __consumer_offsets).
  • Forgetting the topics must be compacted, or using replication factor 1 in production.

context