Blog

Why Kafka Connect Is Not Set-and-Forget Infrastructure

Table of Contents

Table of Contents

A connector can still be listed as deployed while the business flow behind it has stopped. Illustrative scenario, not a customer case: a sink connector sends records to an external system that begins rejecting requests after a credential change. One task enters a retry loop, another fails, and the worker remains healthy enough to answer its REST calls. The platform dashboard still shows a connector object. The downstream table stops receiving fresh data.

An operator restarts the task. It comes back, reads from the same boundary, and meets the same rejected request. A second restart creates the same outcome. The failure is no longer a one-off task event; it is an operating gap. Someone must decide what to retry, which offset is safe, how to repair the dependency, and whether the worker has enough capacity to catch up afterward.

Kafka Connect is a framework for running integrations, not an exemption from integration operations. Each connector creates three standing responsibilities: failure retry, offset management, and capacity change. If those duties have no owner or measurable signal, a connector can become a restart black hole: it appears active, consumes attention, and makes no useful progress.

1A connector is a running system with state

Kafka Connect gives teams a common runtime for source and sink connectors, tasks, distributed workers, and REST management. That standard model can make a connector look more declarative than it is: a JSON configuration hides a process with a source or destination, credentials, network paths, plugin dependencies, Kafka state, and recovery behavior.

In distributed mode, Connect workers use Kafka topics to store connector configuration, offsets, and task status. The Apache Kafka Connect documentation describes the framework and its distributed operation. Those internal topics are part of the connector's operating state. A worker restart can be routine when that state is healthy; rebuilding or moving the runtime without understanding it can change where a source resumes or how a sink replays records.

Liveness is not progress. A worker can answer health checks while a task is repeatedly failing to deliver a batch. A connector can report RUNNING while its source is producing nothing because an upstream log is unavailable, or while its sink is accepting records too slowly to meet its freshness objective. The useful question is not “is the connector deployed?” It is “what evidence shows that the expected boundary is advancing?”

That question sets the operating model. Every connector needs a known owner, a defined progress signal, a recovery boundary, and a capacity assumption. The details vary by plugin, but the responsibilities do not disappear when the configuration is short.

Connector responsibility map showing retry, offset, and capacity duties around a Kafka Connect task

2The three responsibilities every connector carries

The three responsibilities are connected. A retry policy changes when an offset can move, while a capacity change can trigger task reassignment and alter which worker owns recovery. A connector can therefore look healthy in one dashboard and unhealthy in the data path.

2.1Failure retry is a policy, not a button

A failed request can be transient, terminal, or ambiguous. A destination may return a temporary rate limit, reject a record because its shape is invalid, or accept a write while the response is lost. The connector needs a policy for each class. That policy may include bounded retries, backoff, error tolerance, a dead-letter topic, or a pause for human repair. The right choice depends on the connector and the destination's delivery semantics.

The dangerous default is an unexamined retry loop. Retrying can protect against a short dependency interruption, but it can also keep a task busy while its offset remains fixed. A dead-letter path can isolate malformed records, but it does not define replay ownership or prove that the sink and topic are consistent afterward.

A retry runbook should therefore answer four questions:

  • Which error classes are safe to retry, and which require configuration or data repair?
  • How is backoff bounded so a task cannot occupy a worker forever without progress?
  • When is a record routed to a dead-letter topic, and what metadata travels with it?
  • What signal confirms that recovery moved the connector beyond the failing boundary?

The last question is the one restarts tend to skip. Restarting changes process state. It does not change a bad credential, an invalid record, a closed network path, or a sink that cannot accept the write.

2.2Offset management is a data contract

Offsets make restart and replay possible, but an offset is meaningful only at a defined boundary. For a source connector, the stored position may represent a source log position, a table snapshot boundary, or another plugin-specific coordinate. For a sink connector, Kafka consumer progress describes what the task has processed according to the connector framework. Neither label alone proves that the destination has applied a business transaction in the way the application expects.

Operators need to know where progress is stored, when it is committed, and what happens if the process stops between the external write and the offset update. That determines whether a restart replays records, skips records, or requires destination reconciliation. Idempotent writes can make replay safe for some sinks; others need keys, upserts, deduplication, or a repair job.

Offset ownership also matters during a task rebalance or worker replacement. A source connector may resume from a stored source position while a sink connector may resume from Kafka consumer offsets. A configuration export without the associated state is not a recovery plan. The connector lifecycle readiness checklist treats offset handling as part of lifecycle readiness for this reason.

Write down the recovery boundary in language an operator can use: “resume from the last committed source position,” “replay from the last Kafka offset,” or “reconcile the destination before resuming.” If nobody can state the boundary, the connector is not ready for an unattended restart.

2.3Capacity changes alter behavior

Capacity is more than worker CPU and memory. It includes source parallelism, sink write limits, Kafka partition count, task count, plugin behavior, network bandwidth, destination quotas, and the backlog a connector must clear after an interruption. Increasing tasks.max can create more parallel work, but the plugin may not use it, the source may serialize access, or the sink may throttle the additional requests.

A capacity change can also move task ownership. Workers rebalance, tasks restart, and in-flight work can meet the offset boundary described above. Adding workers without checking source and sink limits can increase contention; adding tasks without a backlog target can improve throughput while leaving freshness outside the objective.

Use a capacity review that starts with the bottleneck, not with a preferred knob:

QuestionEvidence to collectDecision it supports
Is the source producing work fast enough?Source freshness, snapshot state, or source log positionAdd source parallelism or repair the source dependency
Can the sink accept the offered rate?Batch latency, rejected requests, destination quota, and commit progressTune batches, reduce pressure, or expand the destination
Is the Connect runtime constrained?Worker resource use, task restarts, rebalance activity, and queue depthAdd worker capacity or isolate workloads
Can recovery absorb the backlog?Backlog size, expected catch-up rate, and production rateChoose a safe recovery window and protect live traffic

A capacity change is complete only when the connector's progress and freshness signals improve. A larger worker pool that leaves offsets fixed is a configuration change, not recovery.

3Silent task death creates restart black holes

Task failure is visible only when the monitoring surface follows task state. A connector-level status can remain present while an individual task is failed, paused, or repeatedly restarting. The inverse can also happen: the task process is alive, but a dependency error prevents useful progress. Both cases are operationally dangerous because a shallow health check answers the wrong question.

A restart black hole has a recognizable shape. An operator sees a failed or stalled task, restarts it, and watches it return to the same state. The action feels productive because the task ID changes state and the worker emits fresh logs. The data boundary does not move. The next responder repeats the same action because the runbook records “restart connector” without recording the failing dependency or the last successful offset.

Make restart conditional. Capture task state, last known progress, error class, dependency, and offset boundary. After changing the cause, restart once and verify progress. If progress remains flat, stop restarting and escalate with the evidence collected.

The following timeline is useful during an incident because it separates process state from data progress:

Illustrative silent task death and restart black hole timeline for Kafka Connect

  1. Detect: a task is failed, paused, or running without expected progress.
  2. Locate: identify the last successful offset or source position and the dependency that accepted the last useful work.
  3. Classify: separate a retryable dependency fault from a bad record, an authentication problem, a capacity limit, or an ambiguous write.
  4. Repair: change the cause or choose a documented replay, dead-letter, or reconciliation path.
  5. Verify: confirm offset movement, freshness recovery, and destination acceptance before closing the event.

The sequence leaves an audit trail: what was tried, which boundary was protected, and why another restart is or is not safe.

4A minimal connector monitoring set

A minimal set should observe the data path beyond the worker process. It gives an on-call responder enough context to decide whether the connector is progressing, retrying, or waiting for a dependency. Thresholds should come from the connector's business freshness and delivery objective; no universal retry count or lag value fits every integration.

SignalWhat it tells youPage or route when
Connector and task stateWhether the runtime is running, paused, failed, or restartingA task is failed or restarting without an approved change window
Progress positionWhether source positions or Kafka offsets are advancingThe position is unchanged while input or expected work exists
FreshnessWhether the downstream view is within its business objectiveSource-to-sink freshness exceeds the connector's declared objective
Retry and error classWhether errors are transient, terminal, or ambiguousThe same error repeats while progress stays flat, or a terminal class appears
Dead-letter volumeWhether records are leaving the normal delivery pathVolume crosses the owner-defined review threshold
Throughput and backlogWhether the connector can catch up without harming live trafficBacklog grows, catch-up rate falls, or destination throttling dominates
Ownership and dependencyWho can change the failing conditionAny alert lacks an owner, escalation path, or recovery action

Correlate these signals in one alert or runbook view. A RUNNING task with flat offsets and repeated authorization errors needs access repair; a RUNNING task moving records into a dead-letter topic needs a data-quality decision. A single “connector healthy” label cannot distinguish them. The connector fleet observability guide expands the model into source freshness, Kafka movement, sink delivery, and ownership signals.

Connector SLI card linking progress, freshness, delivery, recovery, and capacity to on-call actions

5The on-call card should fit beside the dashboard

A monitoring set becomes useful when a responder can act from it. Keep the on-call card short enough to read during an incident, but specific enough to prevent a reflex restart.

Connector incident card

Trigger: task failure, flat progress, freshness breach, or a retry loop.

Check: connector and task state; last successful offset or source position; last successful sink acceptance; current error class; source and destination reachability; owner and escalation path.

Choose: repair the dependency, pause and correct data, route to a dead-letter path, or perform a bounded restart with a named replay boundary.

Verify: offsets move, freshness recovers, the sink accepts records, and backlog is not growing faster than it can be cleared.

Record: the cause, protected boundary, action taken, and evidence used to close the event.

The card changes the definition of a successful restart. The process must return, but the data path must also advance. That distinction protects operators from closing an incident because a REST endpoint became responsive again.

6Where a Kafka-compatible platform fits

Once the connector responsibilities are explicit, the platform question is easier to place. A Kafka-compatible platform can preserve the protocol boundary used by existing Connect workers and plugins while changing broker-side storage, scaling, and recovery. That may reduce one class of coupling, but it does not transfer connector ownership to the platform.

AutoMQ is a Kafka-compatible cloud-native streaming platform built on a Shared Storage architecture. Its documentation covers compatibility with Apache Kafka and Kafka Connect management. The relevant evaluation question is whether your existing connector plugins, configurations, state handling, and workload tests pass on the target path.

The boundary should stay clear. A platform change can remove a broker-local constraint; it cannot make a failed sink accept a record or make an ambiguous offset safe. Keep the connector test and the platform test in one plan, but prove connector recovery separately from platform capacity and recovery behavior.

7Connectors deserve SLIs too

Return to the illustrative scenario. The worker was reachable, the connector object existed, and restarts produced activity. None of those facts proved that the destination was receiving the expected data. The missing signal was a connector SLI tied to progress and freshness, with a recovery action that protected the offset boundary.

A practical connector SLI set can be small:

  • Progress SLI: expected source positions or Kafka offsets advance when input exists.
  • Freshness SLI: the downstream result stays within its declared business objective.
  • Delivery SLI: accepted records, retries, and dead-letter routing remain within the connector's contract.
  • Recovery SLI: a documented failure can return to progress without an uncontrolled replay or gap.
  • Capacity SLI: the connector has enough worker, source, sink, and backlog headroom for the declared workload.

These SLIs make the three responsibilities visible. Retry health appears as bounded recovery rather than endless attempts. Offset management appears as a known and testable boundary. Capacity changes appear as restored progress and freshness rather than a larger worker count.

Kafka Connect is valuable because it turns many integrations into a common operating surface. That surface deserves the same engineering discipline as the Kafka brokers behind it. If a connector is important enough to move production data, give it an owner, an SLI, and an on-call card. When you are ready to test that operating model on a Kafka-compatible platform, use the AutoMQ BYOC evaluation path with one representative connector and its recovery runbook.

Newsletter

Subscribe for the latest on cloud-native streaming data infrastructure, product launches, technical insights, and efficiency optimizations from the AutoMQ team.

Join developers worldwide who leverage AutoMQ's Apache 2.0 licensed platform to simplify streaming data infra. No spam, just actionable content.

I'm not a robot
reCAPTCHA

Never submit confidential or sensitive data (API keys, passwords, credit card numbers, or personal identification information) through this form.