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

使用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_

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 10:32:28