Command Palette

Search for a command to run...

Hectal
PHASE 9Intermediate ~7 min· topic 1 of 4

Topic 9.1

Kafka Connect: Workers, Connectors, and Tasks

In one line

Kafka Connect is a framework for streaming data between Kafka and other systems using configuration instead of code. Source connectors import (databases, files, APIs), sink connectors export (search, S3, warehouses, databases). Workers run connectors split into tasks, and Connect stores configs, offsets and status in Kafka topics.

0/4 · 0%

Think of it like this

Universal adapters at an airport charging station. Instead of every traveller building their own charger, standard adapters connect any device to the grid. Connectors are those adapters between Kafka and other systems.

Key ideas

  1. 01

    Concepts: a connector (a configured instance of a plugin, like a JDBC sink for one table), tasks (units of parallel work; a sink's tasks divide topic partitions, a source's divide tables or shards), workers (JVM processes running tasks), converters (bytes ↔ Connect records: JSON, Avro, Protobuf), and single message transforms (SMTs, small per-record changes like renaming fields or routing).

  2. 02

    Standalone mode: one worker, config files, offsets on local disk; for development. Distributed mode: many workers form a group; configs, offsets and status live in Kafka topics (connect-configs, connect-offsets, connect-status); tasks rebalance across workers on failure.

  3. 03

    Management is a REST API: POST /connectors with JSON config, GET /connectors/<name>/status, PUT .../config, POST .../restart?includeTasks=true.

  4. 04

    Examples: Debezium PostgreSQL/MySQL sources, JDBC source/sink, S3 sink (Parquet or JSON files by time partition), Elasticsearch/OpenSearch sink, BigQuery or Snowflake sinks.

  5. 05

    Delivery: sink connectors usually give at-least-once (make sinks idempotent via upserts or keys); some support exactly-once, and source connectors can run exactly-once since Kafka 3.3 (KIP-618) when enabled.

Code & diagrams

connect.mermaiddiagram
Rendering diagram…
s3-sink.jsonjson
{
  "name": "orders-s3-sink",
  "config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "tasks.max": "6",
    "topics": "commerce.orders",
    "s3.bucket.name": "acme-lake",
    "s3.region": "ap-south-1",
    "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
    "partitioner.class": "io.confluent.connect.storage.partitioner.TimeBasedPartitioner",
    "path.format": "'dt'=YYYY-MM-dd/'hr'=HH",
    "partition.duration.ms": "3600000",
    "flush.size": "50000",
    "rotate.interval.ms": "600000",
    "locale": "en-IN",
    "timezone": "UTC",
    "errors.tolerance": "all",
    "errors.deadletterqueue.topic.name": "orders-s3-sink-dlq",
    "errors.deadletterqueue.context.headers.enable": "true"
  }
}

Explain it without notes

01

Explain connector, task, worker and converter.

Practice

01

Run a Connect worker in the lab and create a FileStreamSink (or S3-compatible MinIO sink) connector via the REST API.

Trade-offs

  • ↔

    Connect removes custom integration code but adds a cluster to operate and connector-specific quirks to learn.

Done when you can

  • I can explain Connect's model and deploy source and sink connectors.