You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何修改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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.27 09:25:33