如何修改Kafka-Connect创建的Elasticsearch索引名称?
修改Kafka Connect Elasticsearch Sink的自定义索引名称
Elasticsearch Sink Connector默认以Kafka主题名作为索引名称,要自定义索引名,核心是通过Kafka Connect的转换器重写目标索引标识(连接器最终将转换后的topic名作为索引名)。以下针对两种常见需求给出修改方案:
方案1:固定自定义索引名称
如果要将所有数据写入固定名称的索引(例如my-proxy-trace),调整RegexRouter转换器配置即可:
修改后的完整请求命令:
curl --location --request POST 'localhost:8084/connectors' --header 'Content-Type: application/json' --data-raw ' { "name": "PROXY_HTTP_TRACE", "config": { "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector", "tasks.max": "1", "key.ignore": "true", "schema.ignore": "true", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "key.converter.schemas.enable": "false", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "connection.url": "http://localhost:9200", "connection.username": "elastic", "connection.password": "password", "type.name": "_doc", "topics": "PROXY_HTTP_TRACE", "transforms": "renameIndex", "transforms.renameIndex.type": "org.apache.kafka.connect.transforms.RegexRouter", "transforms.renameIndex.regex": ".*", "transforms.renameIndex.replacement": "my-proxy-trace" } }'
关键修改说明:
- 移除原有无效的
dropPrefix转换(原配置中该转换未改变任何内容) - 添加
renameIndex转换器:用.*匹配原主题名,直接替换为你指定的固定索引名
方案2:带时间分片的自定义索引名称
如果需要保留按月分片的逻辑,同时自定义索引前缀(例如custom-proxy-trace-202405),调整TimestampRouter的格式配置即可:
修改后的完整请求命令:
curl --location --request POST 'localhost:8084/connectors' --header 'Content-Type: application/json' --data-raw ' { "name": "PROXY_HTTP_TRACE", "config": { "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector", "tasks.max": "1", "key.ignore": "true", "schema.ignore": "true", "key.converter": "org.apache.kafka.connect.storage.StringConverter", "key.converter.schemas.enable": "false", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "connection.url": "http://localhost:9200", "connection.username": "elastic", "connection.password": "password", "type.name": "_doc", "topics": "PROXY_HTTP_TRACE", "transforms": "routeTS", "transforms.routeTS.type": "org.apache.kafka.connect.transforms.TimestampRouter", "transforms.routeTS.topic.format": "custom-proxy-trace-${timestamp}", "transforms.routeTS.timestamp.format": "YYYYMM" } }'
关键修改说明:
- 移除无效的
dropPrefix转换 - 修改
routeTS的topic.format为自定义前缀+时间变量,最终生成带年月后缀的自定义索引名
注意事项
- 确保目标Elasticsearch中不存在同名的只读索引,否则连接器无法写入数据
- 新配置生效后,会自动创建新的自定义索引,原有旧索引不会被修改或删除
- 若使用多个转换器,需注意
transforms参数中的顺序,转换会按顺序执行
内容的提问来源于stack exchange,提问作者Emrahall
相关产品推荐
相关产品推荐

