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

配置SimplifiedJson后Mongo Kafka Source导致Elasticsearch Sink连接器失效

问题根因

报错核心原因是Elasticsearch的目标索引已存在固化映射:未启用SimplifiedJson配置时,MongoDB输出的$numberLong封装结构被ES自动将对应字段映射为object类型,切换为原生数值输出后类型不匹配,导致解析失败。

解决方案

步骤1:清理无效索引映射

删除原有目标索引,清除已固化的错误类型映射:

DELETE /localdb

可通过Kibana Dev Tools或curl命令执行上述请求。

步骤2:调整连接器配置

Source连接器完整配置

name=mongo-source-connector11
connector.class=com.mongodb.kafka.connect.MongoSourceConnector
tasks.max=1
output.format.value=json
output.json.formatter=com.mongodb.kafka.connect.source.json.formatter.SimplifiedJson

# Connection and source configuration
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false 
connection.uri=mongodb://localhost:27017
database=localdb
collection=abcd
pipeline : [{'$match': { $or: [ { 'operationType': 'insert' }, { 'operationType': 'update' }, { 'operationType': 'delete' }, { 'operationType': 'replace' }, { 'operationType': 'rename' } ] } }]
poll.max.batch.size=1000
poll.await.time.ms=1000

output.schema.infer.value=true

output.format.key=schema
output.schema.key={\"name\":\"ClassroomId\",\"type\":\"record\",\"namespace\":\"com.mongoexchange.avro\",\"fields\":[{\"name\":\"documentKey._id\",\"type\":\"string\"}]}

change.stream.full.document=updateLookup
copy.existing=true

主要调整点:

  • 明确指定output.format.value=json,配合简化JSON格式化器输出标准JSON结构
  • 将value.converter统一为JsonConverter,和Sink端配置对齐

Sink连接器完整配置

name=elasticsearch-sink-connector11
connector.class=io.confluent.connect.elasticsearch.ElasticsearchSinkConnector
tasks.max=1

topics=localdb.abcd
connection.url=http://localhost:9200
topic.index.map=localdb.abcd:localdb
type.name="_doc"
key.converter.schemas.enable=false
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
key.converter=org.apache.kafka.connect.json.JsonConverter
key.ignore=false
schema.ignore=true
behavior.on.malformed.documents=fail

# 提取fullDocument作为写入ES的根文档,避免嵌套层级
transforms=extractDoc
transforms.extractDoc.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.extractDoc.field=fullDocument

主要调整点:

  • 新增schema.ignore=true,关闭schema校验适配无schema的JSON输出
  • 新增SMT转换直接提取变更内容的fullDocument作为写入ES的根文档,简化字段层级
  • 统一value转换器和Source端对齐

步骤3:重启连接器

删除原有已部署的Source和Sink连接器,提交调整后的配置重新启动,触发全量数据同步。新索引会自动将数值类型字段映射为原生long/integer类型,不会再出现$numberLong类封装结构。

特殊场景适配

如果需要保留变更事件的元数据(如操作类型、变更时间等),无需配置SMT转换,提前手动创建ES索引映射,明确指定所有数值字段的类型为long/integer即可避免类型冲突。

内容的提问来源于stack exchange,提问作者java-is-problematic

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 11:45:02