Learn Labs
9. Building Data Pipelines (Kafka Connect)

9.8 A deeper look at Connect internals

Note: "The JSON converter can be configured to either include a schema in the result record or not — so we can support both structured and semistructured data."

8.1 Connectors vs tasks — a clean separation

  1. “Determining how many tasks will run for the connector”
  2. “Deciding how to split the data-copying work between the tasks”
  3. “Getting configurations for the tasks from the workers and passing them along”

The JDBC source example, concretely:

  1. Connect to the database
  2. Discover the existing tables to copy
  3. Decide how many tasks are needed: task_count = MIN( tasks.max , number_of_tables )
  4. Generate a config for each task, using the connector config (e.g. connection.url) and a list of tables it assigns to that task. ► taskConfigs() “returns a list of maps (i.e., a configuration for each task we want to run)”
  5. Workers start the tasks, each with its own unique configuration, “so that it will copy a unique subset of tables from the database”

⚠️ "Note that when you start the connector via the REST API, it may start on ANY node, and subsequently the tasks it starts may ALSO execute on ANY node."

via the workervia the workerexternal systemsource taskpoll → lists of recordsKafkasink taskwrite to the external systemexternal systemTasks are “responsible for actually getting the data in and out”.
  • All tasks are initialized by receiving a context from the worker, then started with a Properties object — the config the connector made.
  • Source context: “an object that allows the source task to store the offsets of source records” — file connector: positions in the file; JDBC source: a timestamp column in a table.
  • Sink context: “methods that allow the connector to control the records it receives from Kafka” — applying back pressure, retrying, and “storing offsets externally for exactly-once delivery”.
  • Source tasks “poll an external system and return lists of records that the worker sends to Kafka brokers”. Sink tasks “receive records from Kafka through the worker and are responsible for writing them to an external system”.
Figure 9.8.3

"Storing offsets externally for exactly-once delivery" is the concrete API behind §2.2's claim — the sink context is what lets a connector commit data + offsets in the target system's transaction.

8.2 Workers — where all the hard operational work lives

Worker responsibilities:

  • Handle the HTTP requests that define connectors and their configuration
  • Store the connector configuration in an internal Kafka topic
  • Start the connectors and their tasks, passing configurations along
  • Failover: “If a worker process is stopped or crashes, other workers in a Connect cluster will recognize that (using the heartbeats in Kafka's consumer protocol) and reassign the connectors and tasks that ran on that worker to the remaining workers.”
  • Scale-out: “If a new worker joins, other workers will notice and assign connectors or tasks to it to make sure load is balanced fairly.”
  • Automatically commit offsets for both source and sink connectors into internal Kafka topics
  • Handle retries when tasks throw errors

"The best way to understand workers is to realize that CONNECTORS AND TASKS are responsible for the 'MOVING DATA' part of data integration, while the WORKERS are responsible for the REST API, CONFIGURATION MANAGEMENT, RELIABILITY, HIGH AVAILABILITY, SCALING, AND LOAD BALANCING."

"This separation of concerns is the main benefit of using the Connect API versus the classic consumer/producer APIs."

Note the reuse: worker failover uses Kafka's consumer group protocol (Ch. 4). Connect didn't invent a membership protocol; a Connect cluster is a consumer group.

8.3 Converters and the Connect data model

Source path
external systemconnector reads an eventConnect data APISchema + Struct/Valuethe configured converterapplied by the workerAvro | JSON | JSON Schema | Protobufprimitive types · byte arrays · stringsKafka

The JDBC source “reads a column and constructs a Connect schema object based on the data types of the columns returned by the database. It then uses the schema to construct a struct that contains all the fields… For each column, we store the column name and the value.”

Sink path — exactly the reverse
Kafka bytesconverterSchema + Valuesink connectortarget
Figure 9.8.58.3 Converters and the Connect data model

"This allows the Connect API to support different types of data stored in Kafka, INDEPENDENT of the connector implementation (i.e., ANY connector can be used with ANY record type, as long as a converter is available)."

Note: "The JSON converter can be configured to either include a schema in the result record or not — so we can support both structured and semistructured data."

8.4 Offset management — the killer feature

"connectors need to know which data they have already processed, and they can use APIs provided by Kafka to maintain information on which events were already processed."

Source connectors: logical partitions and offsets

"the records the connector returns include a LOGICAL PARTITION and a LOGICAL OFFSET. THOSE ARE NOT KAFKA PARTITIONS AND KAFKA OFFSETS but rather partitions and offsets as needed in the SOURCE SYSTEM."

SourceLogical partitionLogical offset
file sourcea filea line number or character number in the file
JDBC sourcea database tablean ID or timestamp of a record in the table

💡 "One of the most important design decisions involved in writing a source connector is deciding on a good way to PARTITION the data in the source system and to TRACK OFFSETS — this will impact THE LEVEL OF PARALLELISM the connector can achieve AND WHETHER IT CAN DELIVER AT-LEAST-ONCE OR EXACTLY-ONCE SEMANTICS."

The write ordering that makes it correct:

  1. Source connector returns records (each with source partition + offset)
  2. Worker Sends the records to kafka brokers
  3. If the brokers successfully acknowledge, the worker Then stores the offsets of the records it sent

⇒ “This allows connectors to start processing events from the most recently stored offset after a restart or a crash.”

Offsets are stored after the data is acked — the same "commit after processing" discipline as Ch. 7 §5.2, applied by the framework so connector authors can't get it wrong.

The three internal topics
offset.storage.topic
source connector logical offsets
config.storage.topic
“the configuration of All the connectors we've created”
status.storage.topic
“the Status of each connector”

⇒ “The storage mechanism is Pluggable and is usually a Kafka topic.”

(Recall Ch. 5 §4.2: these are exactly the config topics that must be 1 partition (strict ordering), 3 replicas (availability), and compacted (indefinite retention).)

Sink connectors: the inverse

"they read Kafka records, which already have a topic, partition, and offset. Then they call the connector put() method that should store those records in the destination system. If the connector reports success, they commit the offsets they've given to the connector back to Kafka, using the usual consumer commit methods."

"Offset tracking provided by the framework itself should make it easier for developers to write connectors and guarantee some level of CONSISTENT BEHAVIOR when using different connectors."


On this page