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
- “Determining how many tasks will run for the connector”
- “Deciding how to split the data-copying work between the tasks”
- “Getting configurations for the tasks from the workers and passing them along”
The JDBC source example, concretely:
- Connect to the database
- Discover the existing tables to copy
- Decide how many tasks are needed:
task_count = MIN( tasks.max , number_of_tables ) - 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)” - 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."
- All tasks are initialized by receiving a context from the worker, then started with a
Propertiesobject — 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”.
"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
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.”
"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."
| Source | Logical partition | Logical offset |
|---|---|---|
| file source | a file | a line number or character number in the file |
| JDBC source | a database table | an 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:
- Source connector returns records (each with source partition + offset)
- Worker Sends the records to kafka brokers
- 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."