使用MongoSourceConnector导入MongoDB数据到Kafka时Decimal字段配置报错,如何正确写入Decimal类型?
解决MongoDB Kafka连接器Decimal128到Avro Decimal的类型转换错误
这个错误的核心原因是MongoDB的Decimal128类型和Avro基于bytes的decimal逻辑类型之间没有被正确转换——连接器读取到的是MongoDB原生的Decimal128对象,但你的Avro schema声明的是bytes类型,两者类型不匹配导致报错。
下面是具体的解决方案,按步骤调整你的配置:
1. 确认数值映射配置的有效性
你已经设置了numeric.mapping='best_fit_or_decimal',这个配置方向是对的:它会让连接器将MongoDB中无法用基本数值类型(long、double)精确表示的数值转换为Decimal128,但需要确保后续Avro转换器能识别这个类型并转换成符合要求的bytes格式。
2. 修正Avro Schema定义
你的当前schema中amount字段的定义虽然符合Avro decimal规范,但需要调整精度/scale匹配业务实际值(比如示例中770.00用scale=2就足够,不需要18),同时修正_id字段的类型错误(示例中_id是字符串,不是long):
修改后的output.schema.value示例:
{ "name":"MongoExchangeSchema", "type":"record", "namespace":"com.mongoexchange.avro", "fields": [ {"name": "dateTime","type": ["long", "null"] }, { "name": "tranProc", "type": [ { "name": "tranProc", "type": "record", "fields": [ {"name": "_id", "type": ["string", "null"]}, {"name": "amount", "type": [{"type": "bytes","logicalType":"decimal", "precision": 10, "scale": 2}, "null"]} ] }, "null" ] } ] }
3. 修正连接器配置的细节
检查并修正Avro转换器的schema registry地址(你的配置中key转换器地址多了一个冗余的.local),然后替换为完整的正确配置:
CREATE SOURCE CONNECTOR mongo_source_tran_stage_st WITH ( 'tasks.max'='1', 'connector.class'='com.mongodb.kafka.connect.MongoSourceConnector', 'output.format.value'='schema', 'output.schema.value'='{ "name":"MongoExchangeSchema", "type":"record", "namespace":"com.mongoexchange.avro", "fields": [ {"name": "dateTime","type": ["long", "null"] }, { "name": "tranProc", "type": [{"name": "tranProc", "type": "record", "fields": [ {"name": "_id", "type": ["string", "null"]}, {"name": "amount", "type": [{"type" :"bytes","logicalType":"decimal", "precision": 10, "scale": 2}, "null"]} ] }, "null" ] } ] }', 'schema.compatibility'= 'BACKWARD', 'change.stream.full.document'='updateLookup', 'pipeline'='[{"$match":{"operationType":{"$in":["insert","update"]}}}]', 'value.converter'='io.confluent.connect.avro.AvroConverter', 'value.converter.schema.registry.url'='http://avro.local:8081', 'key.converter.schemas.enable'=true, 'key.converter'='io.confluent.connect.avro.AvroConverter', 'key.converter.schema.registry.url'='http://avro.local:8081', 'connection.uri'='mongodb://mongo_uri', 'publish.full.document.only'= true, 'topic.prefix'='mongo_struct', 'auto.offset.reset' = 'earliest', 'database'='db', 'collection'='Coll', 'numeric.mapping'='best_fit_or_decimal' );
4. 验证结果
重启连接器后,发送测试数据到MongoDB目标集合,然后:
- 检查Schema Registry中注册的schema是否包含正确的decimal逻辑类型
- 消费Kafka消息,确认
amount字段被正确序列化为符合Avro decimal规范的bytes,且下游系统能正常解析为Decimal类型
如果仍有问题,建议升级MongoDB Kafka Connector到最新稳定版本,确保和Confluent平台的版本兼容性。
内容的提问来源于stack exchange,提问作者roman_
相关产品推荐
相关产品推荐

