02 / 05

What is the difference between source and sink connectors?

Difficulty: 3/10
Connectors, CDC

Source vs Sink Connectors: Data Direction and Interface

The difference is direction: a source connector reads data from an external system and writes it into Kafka; a sink connector reads data from Kafka and writes it to an external system. This directional difference drives almost everything else about how they work. A source connector uses a pull interface: the framework calls poll() on the SourceTask, and the task returns a batch of records it has read from the external system. A sink connector uses a push interface: the framework calls put() on the SinkTask with a batch of records that were read from Kafka. This asymmetry is fundamental. Source tasks must be able to seek to a position in the source system to resume after a failure, which is why source connectors provide source offsets. Sink tasks consume from Kafka partitions and must handle offset commits through the framework, which allows them to recover from failures by resuming from the last committed offset.

The mechanism matters because it determines what each type of connector is responsible for. A source connector must create schemas for the data it emits, because the external system may not have a schema that Kafka understands. It must also track the position in the source system so that it does not re-read or skip data. A sink connector does not create schemas; it receives data with schemas already attached and must validate that the schemas match what the destination system expects. If the schema does not match, the sink connector should throw an exception to signal the error. A source connector applies SMTs before converting data to Kafka's format; a sink connector applies SMTs after converting data from Kafka's format. This ordering is important: a source SMT operates on the connector's internal representation of the external data, while a sink SMT operates on the deserialized Kafka record before it is written to the destination.

A common mistake is to assume that a source connector and a sink connector for the same system are symmetrical. They are not. The JDBC source connector and JDBC sink connector have completely different implementations because reading from a database and writing to a database have different requirements for batching, transactions, and error handling. Another mistake is to ignore the offset management differences. Source connectors need to provide meaningful offsets that allow resumption from the exact position; sink connectors rely on Kafka's offset commit mechanism, and the framework handles it automatically. The trade-off between using a connector and writing a custom application is the same for both directions: connectors give you standardization and operational simplicity, while custom code gives you control. Version note: the SourceTask and SinkTask interfaces have been stable for many versions; the main evolution has been in the quality and number of available connectors, especially Debezium for CDC source connectors and Confluent's sink connectors for various systems.

javascript
  1. 1

    Source connectors read from external systems and write to Kafka; sink connectors read from Kafka and write to external systems.

  2. 2

    Source tasks use a pull interface (poll); sink tasks use a push interface (put).

  3. 3

    Source connectors create schemas and provide source offsets for resumption.

  4. 4

    Sink connectors validate schemas and rely on Kafka's offset commit mechanism.

  5. 5

    Source SMTs apply before conversion to Kafka format; sink SMTs apply after conversion from Kafka format.

  6. 6

    Source and sink connectors for the same system are not symmetrical; they have different implementations.

  7. 7

    Connector interfaces have been stable; the ecosystem of available connectors has grown significantly.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.