配置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
相关产品推荐
相关产品推荐

