Part 7 · 1 chapters · ~8 min

Kafka Connect and CDC

Connect workers, source and sink connectors, distributed mode and task scaling, Debezium for Postgres and MySQL with snapshots and streaming, the event envelope with before and after, the outbox event router, single message transforms, sinks to search and warehouses, and failure handling.

9

Connectors instead of code

code
// a Debezium Postgres source connector (POST to the Connect REST API)
{ "name": "ledger-cdc", "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "db", "database.dbname": "ledger", "database.user": "cdc",
    "plugin.name": "pgoutput", "slot.name": "ledger_cdc", "publication.name": "ledger_pub",
    "table.include.list": "public.outbox",
    "transforms": "outbox", "transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
    "topic.prefix": "ledger" } }

Watch the replication slot: if Connect stops, the slot holds WAL and the database disk fills (Postgres course P8). Alert on slot lag and connector task failures.

KAFKA CONNECT AND CDC
moving data in and out of Kafka without writing consumers by hand
PostgresWALDebezium source connectorlogical slotKafkacdc.public.transfersElasticsearch sinkS3 / warehouse sinkSMTssingle message transforms
swipe the figure sideways, or tap expand for full screen
1/5
Connect
Kafka Connect runs connectors in a cluster of workers: source connectors bring data in, sink connectors write it out, with offsets, retries and scaling handled by the framework.
source and sink connectorsno custom consumer code