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

Kafka Connect MongoDB CDC大消息报错及DLQ配置无效问题求助

问题分析与解决方案

你遇到的RecordTooLargeException属于Kafka Connect框架(Worker)层面的错误,而非MongoDB CDC连接器内部的错误,所以你配置的mongo.errors.*参数无法捕获并转发这类错误,这也是DLQ为空且连接器崩溃的核心原因。以下是具体解决步骤:

一、解决大消息触发的核心错误

调整Kafka端到端的消息大小配置,确保参数匹配:

  • Kafka Broker节点:修改server.properties,设置message.max.bytes=2097152(示例值为2MB,需大于你的实际消息大小)
  • Kafka Connect Worker节点:修改Connect配置文件(如connect-distributed.properties),设置max.request.size=2097152,同时可调整batch.size避免批量消息后超出限制
  • 目标Topic:若Topic有单独配置,需同步调整max.message.bytes,命令示例:
    kafka-topics.sh --alter --topic topic_name --zookeeper zk_host:port --config max.message.bytes=2097152
    

二、配置框架级DLQ捕获发送失败消息

要将这类发送失败的消息转发到DLQ,需配置Kafka Connect框架的全局错误处理参数,而非Mongo连接器的参数:
在Connect Worker配置文件中添加以下内容:

# 开启框架级错误容忍
errors.tolerance=all
# 指定DLQ主题
errors.deadletterqueue.topic.name=topic_dlq
# 可选:设置DLQ消息过期时间(示例为1天)
errors.deadletterqueue.retry.timeout.ms=86400000
# 可选:记录错误详情到DLQ消息头
errors.log.enable=true
errors.log.include.messages=true

分布式模式下需重启所有Connect Worker节点,独立模式重启Connect服务即可生效。

三、修复连接器崩溃问题

日志中的OffsetStorageWriter is already flushing错误,是发送失败后触发的偏移量提交冲突导致的。解决大消息问题或配置框架级错误容忍后,错误会被正确处理,不会触发重复的偏移量提交操作,连接器崩溃问题会自动消失。

内容的提问来源于stack exchange,提问作者Devidb

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 03:10:48