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

MongoDB Source Connector生成MongoID格式Kafka Key及连接器配置故障排查

解决MongoDB Source Connector与Elasticsearch Sink Connector配置问题

问题背景

使用com.mongodb.kafka.connect.MongoSourceConnector和io.confluent.connect.elasticsearch.ElasticsearchSinkConnector搭建数据同步链路,当前配置如下:

MongoDB Source Connector 配置

{
    "name": "ais-mongodb-source",
    "config": {
        "connector.class": "com.mongodb.kafka.connect.MongoSourceConnector",
        "publish.full.document.only": "true",
        "database": "ais-user",
        "output.json.formatter": "com.mongodb.kafka.connect.source.json.formatter.SimplifiedJson",
        "offset.partition.name": "ais-user.1",
        "output.format.value": "json",
        "tasks.max": "1",
        "connection.uri": "",
        "value.converter.schemas.enable": "false",
        "value.converter": "org.apache.kafka.connect.storage.StringConverter",
        "key.converter": "org.apache.kafka.connect.storage.StringConverter",
        "change.stream.full.document": "updateLookup"
    }
}

Elasticsearch Sink Connector 配置

{
    "name": "ais-es-sink-connector",
    "config": {
        "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
        "type.name": "_doc",
        "topics": "ais-user.administrator",
        "tasks.max": "1",
        "key.ignore": "true",
        "schema.ignore": "true",
        "key.converter.schemas.enable": "false",
        "name": "ais-es-sink-connector",
        "value.converter.schemas.enable": "false",
        "connection.url": "http://es01:9200",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "key.converter": "org.apache.kafka.connect.json.JsonConverter"
    }
}

遇到的错误

  1. ES索引失败:
Caused by: org.apache.kafka.connect.errors.ConnectException: Indexing record failed -> Response status: BAD_REQUEST,
 Index: ais-user.administrator,
 Document Id: ais-user.administrator+0+12
  1. 尝试用Transform生成MongoID格式的Kafka Key时,出现Only MAP Supported错误。

错误分析

  1. ES索引失败:key.ignore=true让ES自动生成包含topic、分区、偏移量的文档ID,这类ID可能触发ES的字段验证规则(如特殊字符限制),或在重复同步时引发数据冲突。
  2. Only MAP Supported错误:MongoDB Source使用StringConverter将消息值序列化为JSON字符串,而非结构化的MAP对象,Transform无法直接从字符串中提取_id字段。

解决方案

1. 修改MongoDB Source Connector配置,提取MongoDB的_id作为Kafka Key

将消息值转换器改为JsonConverter,让消息以结构化MAP传递,同时添加Transform提取_id作为Key:

{
    "name": "ais-mongodb-source",
    "config": {
        "connector.class": "com.mongodb.kafka.connect.MongoSourceConnector",
        "publish.full.document.only": "true",
        "database": "ais-user",
        "output.json.formatter": "com.mongodb.kafka.connect.source.json.formatter.SimplifiedJson",
        "offset.partition.name": "ais-user.1",
        "output.format.value": "json",
        "tasks.max": "1",
        "connection.uri": "",
        "value.converter.schemas.enable": "false",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter", // 替换为JsonConverter
        "key.converter": "org.apache.kafka.connect.json.JsonConverter", // 同步修改Key转换器
        "change.stream.full.document": "updateLookup",
        // 添加Transform提取_id作为Kafka Key
        "transforms": "extractKey,extractId",
        "transforms.extractKey.type": "org.apache.kafka.connect.transforms.ValueToKey",
        "transforms.extractKey.fields": "_id",
        "transforms.extractId.type": "org.apache.kafka.connect.transforms.ExtractField$Key",
        "transforms.extractId.field": "_id"
    }
}

2. 修改Elasticsearch Sink Connector配置,使用Kafka Key作为ES文档ID

关闭key.ignore,让ES使用提取的MongoDB _id作为文档ID,避免自动生成的ID问题:

{
    "name": "ais-es-sink-connector",
    "config": {
        "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
        "type.name": "_doc",
        "topics": "ais-user.administrator",
        "tasks.max": "1",
        "key.ignore": "false", // 修改为false
        "schema.ignore": "true",
        "key.converter.schemas.enable": "false",
        "name": "ais-es-sink-connector",
        "value.converter.schemas.enable": "false",
        "connection.url": "http://es01:9200",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "key.converter": "org.apache.kafka.connect.json.JsonConverter"
    }
}

关键说明

  • 使用JsonConverter替代StringConverter,确保消息以结构化MAP形式传递,Transform才能正常提取字段。
  • 提取MongoDB的_id作为Kafka Key,既解决了ES文档ID的格式问题,又能保证数据更新时的幂等性(相同_id会覆盖原有文档)。

内容的提问来源于stack exchange,提问作者Difa Al

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 17:33:16