如何在Confluent 4的Kafka Elasticsearch Sink中自定义索引名称
解决Confluent 4.x Kafka Elasticsearch Sink连接器自定义索引名的问题
嘿,我明白你的需求了——你想让Elasticsearch的索引名和Kafka主题名解绑,不用默认的主题名作为索引名。咱们一步步来搞定这个问题:
1. 先修正配置里的明显错误
看你提供的配置,有个小问题可能导致index.name没生效:
你的type.name配置写的是:
"type.name": "type.name=kafka-connect"
这是不对的!type.name的值应该直接是你想要的类型名称,不需要带上参数名。正确写法应该是:
"type.name": "kafka-connect"
这个错误可能会让连接器抛出异常,进而影响其他配置项的正常工作,包括你设置的index.name。
2. 确认index.name的正确用法
在Confluent 4.x版本的Elasticsearch Sink Connector中,index.name参数就是用来指定固定的Elasticsearch索引名,替代默认的Kafka主题名。你已经配置了"index.name": "asimtest",修正type.name之后,连接器应该就会把数据写入asimtest索引,而不是mysql-foobar了。
3. 修正后的完整配置
这里是调整好的完整配置,你可以直接用:
{ "name": "es-sink-mysql-foobar-02", "config": { "_comment": "-- standard converter stuff -- this can actually go in the worker config globally --", "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector", "value.converter": "io.confluent.connect.avro.AvroConverter", "key.converter": "io.confluent.connect.avro.AvroConverter", "key.converter.schema.registry.url": "http://localhost:8081", "value.converter.schema.registry.url": "http://localhost:8081", "_comment": "--- Elasticsearch-specific config ---", "_comment": "Elasticsearch server address", "connection.url": "http://localhost:9200", "_comment": "Elasticsearch mapping name. Gets created automatically if doesn't exist ", "type.name": "kafka-connect", "_comment": "Custom Elasticsearch index name (instead of using Kafka topic name)", "index.name": "asimtest", "_comment": "Which topic to stream data from into Elasticsearch", "topics": "mysql-foobar", "_comment": "If the Kafka message doesn't have a key (as is the case with JDBC source) you need to specify key.ignore=true. If you don't, you'll get an error from the Connect task: 'ConnectException: Key is used as document id and can not be null.", "key.ignore": "true" } }
4. 额外:更灵活的索引命名方案(如果后续需要)
如果你之后需要基于Kafka主题名做索引名转换(比如把多个主题映射到不同索引,或者对主题名做正则替换),可以用Confluent的Single Message Transforms(SMT)里的RegexRouter。比如想把主题mysql-foobar转换成my-es-index-foobar,可以加这些配置:
"transforms": "renameIndex", "transforms.renameIndex.type": "org.apache.kafka.connect.transforms.RegexRouter", "transforms.renameIndex.regex": "mysql-(.*)", "transforms.renameIndex.replacement": "my-es-index-$1"
不过对你现在的需求来说,直接用index.name就足够简单啦。
最后记得更新连接器配置并重启,你可以用REST API来操作:
curl -X PUT -H "Content-Type: application/json" --data @your-config.json http://localhost:8083/connectors/es-sink-mysql-foobar-02/config
内容的提问来源于stack exchange,提问作者Asim Zaidi
相关产品推荐
相关产品推荐

