Kafka Connect Elasticsearch连接器启动失败:MAP不支持作为文档ID
解决Kafka Connect Elasticsearch启动报错:MAP is not supported as the document id
嘿,我来帮你搞定这个连接器启动失败的问题!这个报错的根源很明确:你用的Elasticsearch连接器默认会把Kafka消息的Key当作Elasticsearch文档的_id,但你的Key是一个JSON对象(也就是MAP类型),而连接器不支持把这种复杂结构直接作为文档ID——文档ID得是字符串、数字这类简单类型才行。
下面给你几个实用的解决方案,你可以根据自己的需求选:
方案1:从Key里挑一个字段当文档ID
看你的示例Key里有个"key":"OKOK"字段,这个看起来很适合当唯一ID。那你可以在连接器配置里这么设置:
name=elasticsearch-sink connector.class=io.confluent.connect.elasticsearch.ElasticsearchSinkConnector tasks.max=1 topics=你的主题名 connection.url=http://localhost:9200 key.ignore=false key.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false document.id.strategy=io.confluent.connect.elasticsearch.FieldValueStrategy document.id.field=key
核心就是document.id.strategy指定用字段值当ID,document.id.field明确取Key里的key字段。
方案2:让Elasticsearch自动生成ID
要是你不需要用Kafka的Key当文档ID,那最简单的办法就是直接忽略Key,让ES自己生成唯一ID。配置改成这样:
name=elasticsearch-sink connector.class=io.confluent.connect.elasticsearch.ElasticsearchSinkConnector tasks.max=1 topics=你的主题名 connection.url=http://localhost:9200 key.ignore=true value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false
只需要把key.ignore设为true就行,省心又快捷。
方案3:把Key转换成简单类型
如果你非得用整个Key的内容当ID,那可以先通过Kafka的Transform功能把Key转成字符串。比如用ExtractField$Key提取Key里的唯一字段,或者把整个JSON序列化成字符串。举个提取字段的配置例子:
name=elasticsearch-sink connector.class=io.confluent.connect.elasticsearch.ElasticsearchSinkConnector tasks.max=1 topics=你的主题名 connection.url=http://localhost:9200 key.ignore=false key.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=false transforms=extractKey transforms.extractKey.type=org.apache.kafka.connect.transforms.ExtractField$Key transforms.extractKey.field=key
这个Transform会把原来的JSON Key转换成只有OKOK这样的字符串,就能正常当文档ID了。
最后验证一下
改完配置重启连接器后,建议用Kafka控制台消费者看看消息格式对不对:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic 你的主题名 --property print.key=true --property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer --property value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
要是Key已经是简单类型了,那连接器应该就能正常启动工作啦!
内容的提问来源于stack exchange,提问作者ardlema
相关产品推荐
相关产品推荐

