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
相关产品推荐
相关产品推荐

