Learn Labs
9. Building Data Pipelines (Kafka Connect)

9.11 Deploy / monitor / scale / recover

Deployment checklist

  • Connect workers on Separate servers from brokers
  • ≥ 2–3 workers for high availability (“at least two or three”)
  • Same group.id across workers = one Connect cluster
  • bootstrap.servers: at least three brokers
  • plugin.path with One subdirectory per connector (jar + deps)
  • Never add connectors to the classpath
  • key.converter / value.converter chosen deliberately (+ schemas.enable or schema.registry.url)
  • Internal topics: offset.storage.topic, config.storage.topic, status.storage.topic — 1 partition, RF 3, Compacted
  • Credentials via an External secret provider, not config files
  • error.tolerance + dead letter queue on sink connectors
  • Standalone mode Only for machine-pinned connectors (syslog)
  • Debezium instead of JDBC polling wherever a CDC connector exists

Scaling levers

What you wantThe lever
More throughput per connectortasks.max — bounded by the connector's own splitting logic, e.g. table count
More capacityAdd workers; the cluster rebalances connectors and tasks automatically
Parallelism on one machineConnect “focuses on parallelizing the work… on a single node as well as by scaling out”
Absorb producer/consumer skewKafka itself is the buffer; scale each side independently
Reduce network/storageCompression

Monitoring

SignalWhy
Connector / task status (status.storage.topic, REST API)Tasks fail individually and stay failed until restarted
Connect worker logs"Missing configurations or libraries are common causes for errors" — the first place to look
Source connector lag (source system position vs stored offset)Whether the source is keeping up
Sink connector consumer lagSinks are consumers — Ch. 7's lag monitoring applies directly
Dead letter queue depthSilent error.tolerance drops are invisible otherwise
Worker rebalance frequencyUses the consumer group protocol → same flapping concerns as Ch. 4
Target system write errorsThe sink's actual job

Recovery / "backup" in pipeline terms

replayidempotent writesall state lives hereKafka retentionthe replay buffersink connectorits consumer group offsetstarget systeminternal compacted topicsoffsets · configs · statusA worker is disposable; losing one loses nothing.
  • Kafka's retention is the pipeline's replay buffer: “it is possible to go back in time and recover from errors when needed… allows replaying the events stored in Kafka to the target system if they were lost.”
  • To rewind a sink: reset the sink connector's consumer group offsets (Ch. 5 §6.4 — stop it first!) and let it re-deliver. Idempotent target writes (deterministic keys) make this safe.
Figure 9.11.2Recovery / 'backup' in pipeline terms

The chapter's closing charge

"Whatever data integration solution you eventually land on, THE MOST IMPORTANT FEATURE WILL ALWAYS BE ITS ABILITY TO DELIVER ALL MESSAGES UNDER ALL FAILURE CONDITIONS. We believe that Kafka Connect is extremely reliable — based on its integration with Kafka's tried-and-true reliability features — but it is important that you TEST the system of your choice, just like we do. Make sure your data integration system of choice can survive STOPPED PROCESSES, CRASHED MACHINES, NETWORK DELAYS, and HIGH LOADS without missing a message. After all, at their heart, data integration systems only have ONE JOB — delivering those messages."

"...It isn't enough that Kafka SUPPORTS at-least-once semantics; you must be sure you aren't ACCIDENTALLY CONFIGURING IT IN A WAY that may end up with less than complete reliability."

(That's the same charge as Ch. 7 §6 — validate, don't assume — now applied to the integration layer.)


On this page