How does a sink connector decide which Kafka topics to consume, and how are partitions assigned to its tasks?
answer
- topics OR topics.regex (mutually exclusive)
- Sink tasks = one consumer group
- 1 partition -> 1 task
- Parallelism capped by total partitions
- open()/close() on assignment change
basics
~20 sA sink connector subscribes to topics via the topics config (explicit list) or topics.regex (pattern). The tasks form a consumer group, so Kafka's group rebalance assigns each partition of those topics to exactly one task.
solid answer
~50 sSink connectors are required to set exactly one of two mutually-exclusive properties: topics (a comma-separated explicit list) or topics.regex (a regular expression matching topic names). Connect builds a Kafka consumer subscription from that. Because all tasks of a sink connector belong to one consumer group (group.id = connect-<connector-name>), the standard consumer rebalance protocol assigns every partition of the matched topics across the tasks — each partition goes to exactly one task at a time. This is why effective sink parallelism is capped at the total partition count: with 12 partitions and 4 tasks, each task gets ~3 partitions; with 12 partitions and 20 tasks, 8 tasks get nothing. New topics matching topics.regex are picked up on metadata refresh, triggering a rebalance. In SinkTask, open(partitions)/close(partitions) callbacks fire on assignment changes so the task can manage per-partition resources.
go deeper
Know a sink picks topics via topics or topics.regex.
Explain the consumer-group rebalance and one-partition-per-task assignment.
Reason about parallelism caps, regex topic discovery, and open()/close() lifecycle.
Plan partitioning strategy across the pipeline and distinguish input selection from SMT-based output routing.
## Choosing topics: topics vs topics.regex Every sink connector MUST specify which topics to read. There are two ways, and they are **mutually exclusive** (setting both is a config error): - **`topics`** — an explicit comma-separated list: `topics=orders,payments`. Predictable, recommended for stable topologies. - **`topics.regex`** — a Java regex matching topic names: `topics.regex=app\\..*`. Useful when topics are created dynamically; newly created matching topics are discovered on the consumer's metadata refresh and pulled in automatically. ## The consumer group under the hood A sink connector is implemented as a managed **Kafka consumer group**. All of the connector's tasks share one group id (conventionally `connect-<connectorName>`). Connect subscribes that group to the chosen topics. Kafka's **group coordinator** then runs the rebalance protocol and assigns partitions: - **Each topic partition is assigned to exactly one task** at any moment. - The assignor (e.g. cooperative-sticky) decides which task gets which partitions. - Adding/removing tasks, or topics gaining partitions, triggers a rebalance. ## Why partition count bounds parallelism Because one partition -> one consumer/task, the maximum number of *useful* sink tasks equals the **total partition count across all subscribed topics**. If `tasks.max` exceeds that, the surplus tasks are assigned no partitions and idle. Conversely, to scale a sink out you must ensure the source topics have enough partitions. ## SinkTask lifecycle callbacks for assignment When the rebalance changes a task's assignment, Connect invokes: - `open(Collection<TopicPartition>)` — partitions newly assigned; allocate per-partition state (file handles, buffers). - `close(Collection<TopicPartition>)` — partitions revoked; flush and release that state. - `put()` then delivers records only for currently-owned partitions. The `SinkRecord` carries `topic()`, `kafkaPartition()`, and `kafkaOffset()` so the task knows the origin of each record — essential for partition-aware routing (e.g. writing one file per topic-partition in an S3 sink). ## Routing on the way out While topics/topics.regex select the *input*, where records land in the external system is connector-specific and often shaped by SMTs (e.g. `RegexRouter` to rewrite the destination name, `TimestampRouter` for time-bucketed paths). Don't confuse input topic selection (topics/topics.regex) with output routing (SMTs + connector config). ## Source connectors differ Source connectors do NOT use topics/topics.regex for input — they generate data. They decide destination topics in their own config (e.g. a topic prefix) or per-record via `SourceRecord.topic()`. topics/topics.regex are a sink-only concept.
- Can you set both topics and topics.regex on the same sink connector?No — they are mutually exclusive; specifying both is a configuration validation error. Pick one.
- Do source connectors use topics.regex to choose input?No. topics/topics.regex are sink-only (they select Kafka topics to consume). Source connectors generate data and decide destination topics via their own config or per SourceRecord.topic().
saying these in an interview costs you the question
- Saying a single partition can be assigned to multiple sink tasks simultaneously.
- Using topics and topics.regex together.
- Claiming source connectors subscribe via topics.regex.
- Thinking more tasks than partitions increases throughput.