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.idacross workers = one Connect cluster bootstrap.servers: at least three brokersplugin.pathwith One subdirectory per connector (jar + deps)- Never add connectors to the classpath
key.converter/value.converterchosen deliberately (+schemas.enableorschema.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 want | The lever |
|---|---|
| More throughput per connector | tasks.max — bounded by the connector's own splitting logic, e.g. table count |
| More capacity | Add workers; the cluster rebalances connectors and tasks automatically |
| Parallelism on one machine | Connect “focuses on parallelizing the work… on a single node as well as by scaling out” |
| Absorb producer/consumer skew | Kafka itself is the buffer; scale each side independently |
| Reduce network/storage | Compression |
Monitoring
| Signal | Why |
|---|---|
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 lag | Sinks are consumers — Ch. 7's lag monitoring applies directly |
| Dead letter queue depth | Silent error.tolerance drops are invisible otherwise |
| Worker rebalance frequency | Uses the consumer group protocol → same flapping concerns as Ch. 4 |
| Target system write errors | The sink's actual job |
Recovery / "backup" in pipeline terms
- 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.
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.)