Kafka Connect需抛出何种错误才能将异常消息路由至DLQ
我有个可能比较基础的问题,目前实在有点困惑。
我之前读过一篇讲解非常全面的、关于Kafka Connect如何处理死信队列(DLQ)的文章,已经参考思路在自己的连接器上完成了如下配置:
errors.tolerance: all errors.log.enable: true errors.log.include.messages: false errors.deadletterqueue.topic.name: com.kafkaconnect.jdbc.dlq errors.deadletterqueue.topic.replication.factor: 1 errors.deadletterqueue.context.headers.enable: true
但实际运行时出现了两个不符合预期的问题:
- 配置的DLQ主题始终为空
- 只要处理的主题中出现一条错误消息,整个处理流程就会停止,连接器直接进入错误状态
我想确认两个核心问题:
- 我原本认为,在
ValueConverter或KeyConverter处理阶段失败的消息会被送入DLQ,这个预期是否正确? ValueConverter处理阶段需要抛出什么类型的错误,才能让任务/连接器继续处理其余消息,同时将异常消息送入DLQ?
目前我在自定义ValueConverter中处理失败逻辑的代码如下,这段代码运行时会直接导致连接器崩溃,而不是将异常消息送入DLQ:
public SchemaAndValue toConnectData(String topic, Headers headers, byte[] value) { try { Object deserialized = this.deserializer.deserialize(topic, value); if (deserialized == null) { return SchemaAndValue.NULL; } else if (deserialized instanceof Message) { Message message = (Message) deserialized; return this.protobufData.toConnectData(message.getDescriptorForType(), message, topic); } throw new DataException(String.format( "Unsupported type returned during deserialization of topic %s ", topic )); } catch (SerializationException e) { throw new DataException(String.format( "Failed to deserialize data for topic %s to Protobuf: ", topic ), e); } }
我知道目前提供的内容不是最小可复现问题场景,但确实不知道如何简单搭建复现环境,想请教有Kafka Connect使用经验的开发者,通常要如何实现才能让反序列化失败的消息正常进入DLQ。
问题根因
首先纠正一个认知偏差:Kafka Connect 3.2之前版本的默认错误处理逻辑(含DLQ路由),仅覆盖连接器自身put/poll阶段、单消息转换(SMT)阶段抛出的DataException,Converter阶段抛出的异常会被直接判定为任务级致命错误,触发任务失败,根本不会进入DLQ路由流程。
这就是你配置了errors.tolerance=all依然会出现任务崩溃、DLQ为空的核心原因。
可行解决方案
方案1:高版本原生支持(推荐)
如果你的Kafka Connect版本≥3.2,直接在现有配置基础上追加两个参数,就能让Converter阶段的错误被DLQ逻辑正常捕获:
# 开启Converter阶段的错误容忍 errors.tolerance.converter = all # 允许Converter阶段的失败消息路由到DLQ errors.deadletterqueue.converter.errors.enable = true
配置生效后,你现有代码里抛出的DataException会被框架正常拦截,失败消息写入DLQ的同时,任务会继续处理后续消息,不会中断。
方案2:低版本兼容处理
如果使用的是3.2以下的版本,没有上述配置项,就不能在Converter内部直接抛异常。你需要把反序列化失败的消息包装成带错误标记的特殊SchemaAndValue对象返回,在后续的SMT或者连接器处理逻辑中识别这个错误标记,再手动抛出DataException——只有Converter之外的阶段抛出的DataException,才能被旧版本的错误处理逻辑捕获并路由到DLQ。
避坑提示
- Converter内抛出的异常必须是
org.apache.kafka.connect.errors.DataException类型,如果抛出其他异常(比如未捕获的SerializationException、RuntimeException),哪怕开了高版本的Converter错误容忍,也会被判定为致命错误导致任务崩溃 - 提前确认DLQ对应主题已经创建,或者Connect集群开启了主题自动创建权限,否则DLQ写入失败也会触发任务报错
- 排查问题时可以直接搜索连接器日志中
DeadLetterQueueReporter关键字,能看到DLQ写入失败的具体原因
内容的提问来源于stack exchange,提问作者Miguel Costa

