You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka Connect需抛出何种错误才能将异常消息路由至DLQ

Kafka Connect配置死信队列不生效,反序列化失败直接导致连接器崩溃

我有个可能比较基础的问题,目前实在有点困惑。

我之前读过一篇讲解非常全面的、关于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主题始终为空
  • 只要处理的主题中出现一条错误消息,整个处理流程就会停止,连接器直接进入错误状态

我想确认两个核心问题:

  1. 我原本认为,在ValueConverter或KeyConverter处理阶段失败的消息会被送入DLQ,这个预期是否正确?
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 21:24:28