9. Building Data Pipelines (Kafka Connect)
9.10 What actually breaks in production — Ch. 9 consolidated
| # | Symptom | Root cause | Fix |
|---|---|---|---|
| 1 | A tangle of bespoke pipelines nobody can maintain; adopting a new system takes a quarter | Ad hoc pipelines — one custom tool per pair of endpoints | Single integration substrate (Connect); think about the whole graph, not the immediate hop |
| 2 | A DBA adds a column and every downstream app breaks (or all must deploy together) | Loss of metadata — no schema propagation or evolution | Schema Registry + a converter that carries schemas |
| 3 | Downstream team needs a field the pipeline dropped years ago; historical data unrecoverable | ETL over-processing; "historical data will require reprocessing (assuming it is available)" | ELT-lean: only transform what benefits every consumer |
| 4 | Pipeline requires constant changes as downstream needs shift | Extreme processing coupling downstream systems to pipeline-time decisions | Preserve raw data; let Streams apps decide |
| 5 | Credentials leaked from a Connect config file | Connector configs contain credentials in plaintext | External secret config providers (Vault / AWS / Azure) |
| 6 | Bad records discovered days later; can't reprocess | Retention shorter than bug-discovery latency | "Kafka can be configured to store all events for long periods" — size retention to your detection latency |
| 7 | Connect worker OOM / broker page-cache contention | Connect running on the broker machines | "run Connect on SEPARATE SERVERS from your Kafka brokers" |
| 8 | Connector plug-in not found | Dependencies placed at the top level of plugin.path instead of in a per-connector subdirectory | One subdirectory per connector, containing the jar and all its dependencies |
| 9 | Bizarre NoSuchMethodError / version conflicts | Connectors added to the Kafka Connect classpath, bringing a dependency that conflicts with Kafka's | Use plugin.path, "the recommended approach" |
| 10 | JDBC connector fails: "Access denied" / driver not found | Driver missing (doesn't ship with the connector for license reasons), or table permissions | Download the MySQL driver into /opt/connectors/jdbc; check the Connect worker log |
| 11 | Deletes never appear downstream; some updates missing | JDBC polling CDC — scans by timestamp / incrementing PK; "relatively inefficient and at times inaccurate" | Debezium (reads the binlog/WAL directly) |
| 12 | Database load spikes from the pipeline | JDBC connector repeatedly scanning tables | Log-based CDC |
| 13 | Duplicate documents in Elasticsearch after reprocessing | Kafka records had null keys (JDBC doesn't populate them) and the sink generated new IDs | key.ignore=true → deterministic topic+partition+offset document ID |
| 14 | File-based pipeline loses data | FileStream connectors — "many limitations and NO reliability guarantees" | FilePulse / FileSystem Connector / SpoolDir |
| 15 | Corrupt messages halt a sink connector | No error tolerance configured | error.tolerance → silently drop, or route to a dead letter queue |
| 16 | Syslog connector randomly stops receiving data | Connector needs to listen on a specific machine's port, but distributed mode may schedule tasks on any node | Standalone mode for machine-pinned connectors |
| 17 | Only one task runs despite tasks.max=10 | JDBC connector uses MIN(tasks.max, number_of_tables) | Understand the connector's own splitting logic |
| 18 | Connector work distributed unevenly | Connectors and tasks "may start on any node" | Worker rebalancing handles it; inspect via REST API |
| 19 | Source connector reprocesses everything after a crash | Logical offsets not stored / mis-designed | The framework stores offsets after broker ack — a connector authoring bug if it recurs |
| 20 | A single-partition offset/config topic became a bottleneck or lost ordering | Internal topics misconfigured | 1 partition + 3 replicas + compaction (Ch. 5 §4.2) |
| 21 | Consumers can't parse Connect output | JSON converter's schemas.enable mismatch, or Avro registry URL missing | Prefix converter params correctly (key.converter. / value.converter.) |
| 22 | Headers not visible in console consumer | Requires Apache Kafka 2.7+ and --property print.headers=true | Upgrade / add the flag |
| 23 | Sink writes overwhelm the target system | No back pressure | Sink context provides back-pressure methods — a well-written connector uses them |
| 24 | Pipeline can't guarantee exactly-once into a database | Kafka transactions can't span systems (Ch. 8) | Use the sink context's external offset storage — commit data + offsets in the target's transaction |