使用AWS MSK Connect构建MongoDB Sink Connector时ProtoBuf数据写入失败求助
问题解决:MSK Connect MongoDB Sink无法处理ProtoBuf编码数据
问题场景
使用AWS MSK Connect搭建MongoDB Sink Connector时,Kafka Topic中的数据为ProtoBuf编码,配置如下:
connector.class=com.mongodb.kafka.connect.MongoSinkConnector key.converter.schemas.enable=false database=MongoDBMSKDemo tasks.max=1 topics=MongoDBMSKDemo.Source connection.uri= <mongodb connection url> value.converter.schemas.enable=true collection=Sink value.converter=org.apache.kafka.connect.converters.ByteArrayConverter max.batch.size=10 key.converter=org.apache.kafka.connect.storage.StringConverter
运行后数据无法写入MongoDB,报错:
org.apache.kafka.connect.errors.DataException: Could not convert value `[B@51798b7f` into a BsonDocument.
解决方案
报错核心原因是ByteArrayConverter仅将Kafka的字节数据转换为Connect的ByteArray类型,但MongoDB Sink Connector无法直接将原始字节数组解析为BsonDocument。需要添加ProtoBuf反序列化步骤,具体如下:
方法一:使用Confluent ProtoBuf转换器(推荐)
- 替换value.converter
将value.converter改为Confluent的ProtoBuf转换器(需确保MSK Connect环境已安装该插件):value.converter=io.confluent.connect.protobuf.ProtobufConverter - 添加Schema Registry配置
指定Schema Registry地址(用于获取ProtoBuf的schema定义):
配置完成后,转换器会自动将ProtoBuf字节数组反序列化为Connect结构化数据,MongoDB Sink即可直接转换为BsonDocument存入数据库。value.converter.schema.registry.url=<你的Schema Registry地址> value.converter.schemas.enable=true
方法二:自定义转换器(无Schema Registry场景)
若不依赖Schema Registry,可编写自定义Converter实现ProtoBuf反序列化:
- 实现
org.apache.kafka.connect.storage.Converter接口 - 在
fromConnectData方法中,将ProtoBuf字节数组反序列化为Java对象,再转换为BsonDocument或Connect Struct类型 - 将自定义Converter打包为JAR,上传到MSK Connect插件目录
- 在配置中指定自定义Converter:
value.converter=com.yourpackage.YourProtoBufConverter
内容的提问来源于stack exchange,提问作者Sandun Sameera
相关产品推荐
相关产品推荐

