Command Palette

Search for a command to run...

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

Topic 9.2

Running Connect in Production

In one line

Run Connect in distributed mode with its internal topics replicated, pin connector versions, size tasks to partitions, configure error tolerance and connector DLQs, monitor task states and lag, and plan for rebalances when connectors or workers change.

0/4 · 0%

Think of it like this

A fleet of delivery vans shared by several shops. Vans (workers) are assigned routes (tasks); if one breaks down, its routes go to others. The dispatcher's logbook (internal topics) must never be lost.

Key ideas

  1. 01

    Internal topics: config.storage.topic (compacted, 1 partition), offset.storage.topic (compacted, many partitions), status.storage.topic (compacted); all with RF 3. Losing offsets means sources re-read or skip data.

  2. 02

    Error handling: errors.tolerance=none (default, the task fails on the first bad record) vs all (skip bad records) plus errors.deadletterqueue.topic.name for sinks, and errors.log.enable for visibility.

  3. 03

    Scaling: sink tasks ≤ partitions of the source topics; add workers to spread tasks. Tune batch settings per connector.

  4. 04

    Operations: connector and task states (RUNNING, FAILED, PAUSED) must be monitored and alerted; failed tasks don't restart themselves unless you automate POST .../restart. Consumer lag for sink connectors appears under group connect-<connector-name>.

  5. 05

    Isolation: noisy connectors can starve others in a shared cluster; separate Connect clusters per team or criticality, and deploy with Strimzi KafkaConnector resources or GitOps so configs are versioned.

Code & diagrams

connect-ops.shbash
curl -s http://connect:8083/connectors?expand=status | jq '.[] | {name: .status.name, state: .status.connector.state, tasks: [.status.tasks[].state]}'
{"name":"orders-s3-sink","state":"RUNNING","tasks":["RUNNING","RUNNING","FAILED","RUNNING","RUNNING","RUNNING"]}

curl -s http://connect:8083/connectors/orders-s3-sink/tasks/2/status | jq -r .trace | head -3
org.apache.kafka.connect.errors.ConnectException: ... AccessDenied (Service: S3; Status Code: 403)

curl -s -X POST "http://connect:8083/connectors/orders-s3-sink/restart?includeTasks=true&onlyFailed=true"

kafka-consumer-groups.sh --bootstrap-server $B --describe --group connect-orders-s3-sink | head -3

When it breaks

A sink task fails on one bad record with errors.tolerance=none

What you see

The task stops; its partitions' data stops flowing to the sink; nobody notices until the dashboard is empty the next day.

Fix & prevent

Alert on task state; use connector DLQs with errors.tolerance=all for sinks where skipping bad records is acceptable.

Explain it without notes

01

Why must Connect's offset topic be replicated and compacted?

Practice

01

Break a sink connector (wrong credentials), observe the FAILED task, fix the config and restart only failed tasks.

Trade-offs

  • ↔

    Skipping bad records keeps pipelines flowing at the cost of completeness, so pair it with a DLQ and alerts.

Done when you can

  • I can run, monitor and troubleshoot Connect in distributed mode.