MongoDB Kafka Sink Connector大消息截断问题求助:消息超4096字节时序列化失败
触发原因
这个问题的核心不是Kafka Broker或生产/消费者的消息大小限制,而是你使用的JsonConverter(Kafka Connect默认的JSON格式转换器)有一个默认的字节缓冲区大小限制——默认值是4096字节。当消息超过这个大小,转换器会在读取到4096字节时停止解析,导致JSON数据被截断,最终抛出JsonEOFException(JSON格式不完整,提前遇到结束符)。
你之前调整的offset.metadata.max.bytes、max.request.size、message.max.bytes、fetch.max.bytes这些参数,都是控制Kafka集群层面传输消息的最大容量,和Connect组件中负责解析消息的Converter缓冲区完全是两码事,所以修改这些参数无法解决问题。
解决方法
只需要调整Kafka Connect中JSON转换器的缓冲区大小即可,具体操作如下:
针对单个Sink Connector配置:
在你的MongoDB Sink Connector配置文件(或创建Connector的REST请求体)中,添加value.converter.buffer.size参数,设置一个大于你最大消息大小的数值(比如32KB=32768,或者根据实际业务场景调整)。示例配置片段:# 确保指定了JSON转换器 value.converter=org.apache.kafka.connect.json.JsonConverter # 关闭Schema(如果你的消息不需要Schema支持,大多数纯JSON场景都不需要) value.converter.schemas.enable=false # 调整缓冲区大小,这里设置为32KB value.converter.buffer.size=32768全局Connect Worker配置:
如果你希望所有使用JSON Converter的Connector都生效,可以修改Connect Worker的配置文件(比如connect-distributed.properties),找到对应的Converter配置部分,添加或修改buffer.size:key.converter=org.apache.kafka.connect.json.JsonConverter key.converter.buffer.size=32768 value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.buffer.size=32768 value.converter.schemas.enable=false
配置修改完成后,重启Kafka Connect Worker(或者删除并重新创建Sink Connector),之后再发送超过4096字节的JSON消息,就不会再出现截断和序列化异常了。
内容的提问来源于stack exchange,提问作者NineFingers

