Table of Contents
Table of Contents
Three platform teams inherit the same failed connector and produce three different policies. Team A keeps retrying because the target might recover. Team B sends every failed record to a dead-letter topic because the task must keep moving. Team C adds an idempotency check after a duplicate write, but never defines which event identity the check should preserve. Each policy looks reasonable in isolation. Together, they leave operators guessing whether a record is delayed, quarantined, or already applied.
These are illustrative team scenarios, not customer cases. They expose a gap in Apache Kafka® Connect operations: retry, dead-letter handling, and idempotency solve different failure classes. A retry buys time for a transient dependency. A dead-letter queue, or DLQ, preserves a record that needs a separate decision. Idempotency makes a replay safe after a write may already have succeeded. A reliable connector runbook has to compose the three in that order, while keeping their boundaries visible.
The useful rule is compact: set a retry limit, route the records that still cannot be processed to a DLQ, and make every replay safe with a stable idempotency key. The policy must still hold when a task restarts during an unclassified response.
1A connector failure is a classification problem
A connector error is an observation, not yet a policy decision. The worker may know that a call failed, but the operator still needs to classify the failure. Is the destination temporarily unavailable? Is the record malformed for this connector? Did the remote system accept the write before the response disappeared? Those cases need different next actions.
The first team treats all failures as transient, so one poison record can hold every later record. The second treats all failures as permanent, so a growing DLQ can hide a destination outage. The third treats duplicate suppression as error handling, so a replay may be safe while a malformed record loses the evidence needed for repair.
A runbook becomes easier to operate when each mechanism has one job and one explicit limit:
| Mechanism | It is responsible for | It cannot guarantee |
|---|---|---|
| Retry | Waiting through a bounded, recoverable failure | That the remote write did not already succeed |
| DLQ | Preserving a record and failure context for later handling | That the record is correct or will be replayed |
| Idempotency | Making a repeated delivery produce one accepted effect | That the key is correct or retained long enough |
That table also explains why the three mechanisms belong together. Retry handles time. DLQ handles an exception that needs ownership. Idempotency handles uncertainty at the side-effect boundary. Remove any one of them and the remaining two inherit a job they cannot safely perform.
2Retry buys time, not safety
Retry is the right first response when the dependency may recover and the same request is safe to send again. Common examples include a connection reset, a throttling response, a temporary network partition, or a destination that is restarting. The retry window should be finite because a connector is part of a larger pipeline. Holding a partition on one record has a cost in lag, worker capacity, and operator attention.
Use a retry budget that answers three questions before an incident: how long should this task wait, how much delay may grow between attempts, and what happens when the budget expires? The values below are a starting example, not a universal production default. Tune them against the destination’s recovery behavior and the consumer lag that the team can repair.
# Illustrative Kafka Connect error-handling settings
errors.retry.timeout: 30000
errors.retry.delay.max.ms: 5000
errors.log.enable: true
errors.log.include.messages: falseThese settings cover the framework’s record error path. A connector may also have its own destination client retries, batch retries, or request timeout settings. Treat those as a second control layer. If both layers retry without a shared budget, the apparent limit at the worker may conceal a much longer loop inside the connector client.
A retry should preserve the same record identity and enough context to diagnose the attempt. Capture the connector name, task, source topic, Partition, Offset, destination operation, error class, and attempt count. Avoid logging full payloads by default when records contain personal or financial data.
The retry limit is also a handoff point. When it expires, do not silently keep the task blocked and do not treat the record as discarded. Move to the DLQ path with an explicit reason. That boundary is where a retry policy becomes an incident policy.
3A DLQ preserves evidence, not correctness
A DLQ is a durable work queue for records that the normal processing path could not accept. It should carry the original record or a safe representation of it, the source location, the connector and task identity, and the error context needed to decide whether to repair, replay, or reject the record. A topic name alone is not a runbook.
Kafka Connect exposes framework settings for dead-letter handling in its configuration reference. A minimal example looks like this:
errors.tolerance: all
errors.deadletterqueue.topic.name: orders-connector-dlq
errors.deadletterqueue.context.headers.enable: true
errors.log.enable: trueerrors.tolerance: all tells the framework to continue past errors that it can route through this error path. That setting has a sharp edge: a task can keep running while the DLQ quietly accumulates records. Pair it with alerts on DLQ production, task error rate, consumer lag, and the age of the oldest unreviewed DLQ record. If the DLQ itself cannot be produced, that is a separate failure that should page the owner instead of disappearing into a log line.
A DLQ record needs an owner and a disposition: replay unchanged, repair and replay, or hold for a data contract decision. Do not replay the entire topic after fixing one transformation. Select records by stable identity, preserve the source position, and record the operator action.
The DLQ does not make the target write idempotent. A record can arrive in the DLQ after a timeout that followed a successful remote write. Replaying it without the same idempotency key may create a duplicate. That is why the DLQ is the middle step in the policy, not the final safety net.
4Idempotency makes replay safe
Idempotency answers a narrow question: if the same logical event is delivered again, will the target keep one intended effect? It does not mean that a connector never retries or that Kafka delivers a record once. It means the target-side operation can recognize a replay and preserve the intended state.
Start with identity that survives a worker restart and a task reassignment. A producer-assigned event ID is usually stronger than a generated destination ID. If the source has no event ID, a namespaced source position such as cluster/topic/partition/offset can identify one record within one log. That identity has a boundary: a republished event or a backfill with a different offset may represent the same business event while carrying a different position. Write that limitation into the policy.
The target operation must use the key in its atomic write boundary. A key-value sink can overwrite a deterministic key or use an atomic conditional create. A database sink can place a unique idempotency marker and the business mutation in one transaction. An object sink can use a deterministic object name or a committed manifest that tells readers which batch is authoritative. A read followed by a separate write is not a durable deduplication protocol because two tasks can pass the read together.
A minimal database marker makes the boundary concrete:
CREATE TABLE sink_dedupe (
sink_name text NOT NULL,
idempotency_key text NOT NULL,
payload_hash text NOT NULL,
first_seen_at timestamptz NOT NULL,
PRIMARY KEY (sink_name, idempotency_key)
);The first delivery inserts the marker and applies the business change in one transaction. A replay with the same key and payload hash can be acknowledged as already applied. The same key with a different hash is a conflict that needs a controlled error path. It should not be silently swallowed as a normal retry because it may indicate a producer bug, a key collision, or a schema interpretation change.
Idempotency state needs a retention policy. Keep it for the longest replay, backfill, and operator recovery window the team supports. If that window is shorter than source retention, document when replay stops being safe. A DLQ has the same obligation: retention without ownership only postpones the decision.
5Compose the policy: retry cap, DLQ, idempotent replay
The three mechanisms work as a handoff chain. The decision should be deterministic enough that two operators reach the same outcome from the same evidence.
- Classify the failure. Retry only when the dependency may recover and the operation can be repeated with the same event identity.
- Apply the retry budget. Keep the task within a finite timeout and delay ceiling. Measure the attempts and the lag created while the record is waiting.
- Route exhausted or non-retryable records. Publish the record and failure context to a named DLQ. If DLQ publication fails, treat that as an operational fault that needs immediate attention.
- Replay through the target’s idempotent boundary. Preserve the original idempotency key. Apply the target mutation and dedupe decision atomically where the destination supports it.
- Close the loop. Record whether the replay succeeded, was already applied, conflicted, or was rejected. A DLQ is healthy only when its oldest item has a known disposition.
Every error does not need every stage. A schema violation may go directly to the DLQ, while a known duplicate may be acknowledged without another retry. The point is to make exceptions explicit instead of letting defaults decide the business outcome.
| Signal at the failure boundary | Action | Operator evidence |
|---|---|---|
| Timeout, reset, throttling, or temporary unavailability | Retry within the budget | Attempt count, error class, lag, and destination response |
| Malformed record, unsupported mapping, or exhausted retry budget | Publish to the DLQ | Source identity, connector context, error reason, and DLQ publish result |
| Unknown outcome after a remote write | Replay with the same idempotency key | Target conflict or dedupe result, payload hash, and audit record |
| Same key with a changed payload | Stop normal replay and investigate | Key scope, hashes, schema version, and producer owner |
This chain prevents a common category error: using retries to solve poison data, using a DLQ to solve duplicate side effects, or using idempotency to hide an unowned data-quality failure.
6Configurations and signals that survive an incident
A production configuration should show the policy where an operator looks during an incident. Keep framework error settings, connector retries, DLQ ownership, target key derivation, and alert rules versioned together. A change to one layer can invalidate another.
The following example combines the framework controls with comments that point to the runbook. The timeout and delay are placeholders for a tested failure budget. They are not claims about a destination’s recovery time.
# Framework-level record error policy
errors.tolerance: all
errors.retry.timeout: 30000
errors.retry.delay.max.ms: 5000
errors.deadletterqueue.topic.name: orders-connector-dlq
errors.deadletterqueue.context.headers.enable: true
errors.log.enable: true
errors.log.include.messages: false
# Team-owned policy, documented beside the connector config
failure_policy.retryable: "timeouts, connection resets, throttling"
failure_policy.non_retryable: "schema mismatch, invalid key, rejected mapping"
failure_policy.idempotency_key: "event.id or source-cluster/topic/partition/offset"
failure_policy.replay_owner: "streaming-oncall"The team-owned lines are illustrative documentation fields. Keep them in Git, a runbook, or a policy registry. The connector configuration should point to a named decision rather than leave “retry” as an unexplained boolean.
Watch signals from three layers. The connector layer should expose task failures, retry attempts, records sent, records skipped, and restarts. Kafka should expose source lag, DLQ production, DLQ publish failures, and the age of the oldest unreviewed DLQ record. The target should expose request outcomes, conditional-write conflicts, duplicate suppressions, payload-hash conflicts, and the time from replay to final decision. Metric names differ, so map each semantic signal to the deployment’s JMX or exporter name.
An alert should point to an action. A rising retry count asks the operator to check the dependency and the retry budget. DLQ growth asks for a record owner and a disposition. Dedupe hits may be expected during recovery, while payload-hash conflicts require investigation. A dashboard that shows only task health cannot distinguish these states.
7What remains portable when the platform changes?
These policies run at the Kafka Connect and destination boundaries. They are carried by the connector configuration, target contract, key derivation rule, DLQ retention policy, and operational ownership. A Kafka-compatible platform may change broker storage or deployment topology without changing those responsibilities.
AutoMQ is a Kafka-compatible streaming platform that preserves the Kafka-facing protocol and ecosystem surface while using a Shared Storage architecture underneath. Its Kafka compatibility documentation describes the source-side contract. A Connect deployment still needs to test the connector version, converters, task behavior, network path, and destination API for its own workload.
That separation is useful during a migration. The Kafka Connect and AutoMQ pipeline example can frame connectivity testing, while the failure policy remains yours to carry forward. The AutoMQ architecture overview explains the platform layer. For cutover replay, the blue-green Kafka migration guide keeps acknowledged records and rollback in view. None of these replaces the target’s idempotency contract or the connector team’s DLQ ownership.
Treat that portability as an operational test. Export the connector configuration, recreate its DLQ and access controls, replay a controlled record, and verify the same key reaches the same target-side decision. If the migration changes a source identity or topic namespace, update the key scope deliberately. A platform move is a good time to test the policy, not a reason to drop it.
8A Kafka Connect failure policy you can copy
Use the template below as a starting point for each source or sink connector. Fill the blanks with workload-specific decisions and link every value to a test or an owner.
Connector: <name>
Owner: <team or on-call rotation>
Retry when:
<temporary dependency failures and safe-to-repeat operations>
Retry budget:
timeout = <tested duration>
maximum delay = <tested delay>
lag action = <what the on-call does when lag grows>
Dead-letter when:
<non-retryable errors or exhausted retry budget>
DLQ topic:
<topic name and retention owner>
DLQ record must include:
<source topic, Partition, Offset, connector, task, error reason, payload reference>
Replay procedure:
<who selects records, how the source identity is preserved, and how action is recorded>
Idempotency key:
<event ID, namespaced source position, or another durable identity>
Target operation:
<atomic overwrite, conditional write, or transaction with dedupe marker>
Key conflict:
<stop, quarantine, or assigned investigation path>
Dedupe retention:
<replay and backfill window>
Watch:
<retry attempts, task errors, source lag, DLQ growth, target conflicts, dedupe hits>
Test after change:
<failure injection point and expected source and target evidence>
Return to the three teams at the opening. Team A needs a retry ceiling and a handoff. Team B needs DLQ ownership and a replay path. Team C needs a stable identity that protects the target while preserving data-quality evidence. The missing runbook was never a fourth mechanism. It was the composition of the three.
If you are testing this policy on a Kafka-compatible streaming platform, start an AutoMQ evaluation with one controlled retry, one DLQ record, and one replay whose target result you can verify. Keep the policy beside the connector configuration so the next platform change carries the same operational contract.
