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" } }
遇到的错误
- 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
- 尝试用Transform生成MongoID格式的Kafka Key时,出现
Only MAP Supported错误。
错误分析
- ES索引失败:
key.ignore=true让ES自动生成包含topic、分区、偏移量的文档ID,这类ID可能触发ES的字段验证规则(如特殊字符限制),或在重复同步时引发数据冲突。 - 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
相关产品推荐
相关产品推荐

