Skip to content

AxonOps — AI-Native Control Plane for Open Source Data Platforms

Kafka Connect Error Handling

Kafka Connect provides configurable error handling for record processing failures.


{
"errors.tolerance": "none"
}

Connector stops on first error. Use for critical data pipelines where any failure requires attention.

{
"errors.tolerance": "all"
}

Log errors and continue processing. Combine with dead letter queue for failed records.


Route failed records to a DLQ topic for later analysis:

{
"name": "my-connector",
"config": {
"connector.class": "...",
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "dlq-my-connector",
"errors.deadletterqueue.topic.replication.factor": "3",
"errors.deadletterqueue.context.headers.enable": "true"
}
}

When errors.deadletterqueue.context.headers.enable=true:

HeaderDescription
__connect.errors.topicOriginal topic
__connect.errors.partitionOriginal partition
__connect.errors.offsetOriginal offset
__connect.errors.connector.nameConnector name
__connect.errors.task.idTask ID
__connect.errors.exception.class.nameException class
__connect.errors.exception.messageError message
ConnectorTaskSource TopicDLQSinkHeaders include:- Original topic/partition/offset- Exception details- Connector/task infoconsumesuccessfailure + headers

{
"errors.log.enable": "true",
"errors.log.include.messages": "true"
}
ConfigurationDescription
errors.log.enableLog errors to Connect worker log
errors.log.include.messagesInclude record content (may expose sensitive data)

For transient failures:

{
"errors.retry.timeout": "300000",
"errors.retry.delay.max.ms": "60000"
}
ConfigurationDefaultDescription
errors.retry.timeout0Total retry time (0 = no retry)
errors.retry.delay.max.ms60000Max delay between retries

{
"name": "resilient-sink",
"config": {
"connector.class": "...",
"topics": "events",
"errors.tolerance": "all",
"errors.retry.timeout": "300000",
"errors.retry.delay.max.ms": "60000",
"errors.deadletterqueue.topic.name": "dlq-resilient-sink",
"errors.deadletterqueue.topic.replication.factor": "3",
"errors.deadletterqueue.context.headers.enable": "true",
"errors.log.enable": "true",
"errors.log.include.messages": "false"
}
}

// Consume and analyze DLQ
consumer.subscribe(Collections.singletonList("dlq-my-connector"));
while (true) {
ConsumerRecords<byte[], byte[]> records = consumer.poll(Duration.ofSeconds(1));
for (ConsumerRecord<byte[], byte[]> record : records) {
String originalTopic = headerValue(record, "__connect.errors.topic");
String errorClass = headerValue(record, "__connect.errors.exception.class.name");
String errorMessage = headerValue(record, "__connect.errors.exception.message");
log.error("Failed record from {} - {}: {}",
originalTopic, errorClass, errorMessage);
// Analyze and potentially replay
}
}

ScenarioConfiguration
Critical data, no losserrors.tolerance=none
High volume, some loss OKerrors.tolerance=all + DLQ
Transient failures expectederrors.retry.timeout > 0
Debug failureserrors.log.include.messages=true