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

为什么Logstash会把多个Kafka topic的数据写入所有Elasticsearch索引?

问题描述

我搭建了kafka→logstash→elasticsearch的处理管道,在kafka中创建了test-1、test-2两个topic,手动向test-1写入67条消息、向test-2写入67条消息,并编写了两份Logstash配置文件分别对应两个topic的消费和索引输出。
预期结果:Elasticsearch中生成两个独立索引,test-1、test-2各存储67条对应数据
实际结果:运行后test-1和test-2索引均存储了127条数据

相关配置

test1.conf

input {
    kafka {
        bootstrap_servers => "172.29.39.115:9093,172.29.39.116:9093"
        topics => "test-1"
        #consumer_threads => 2
        client_id => "logstash-test-1"
        group_id => "logstash-test-1-0"
        security_protocol => "SSL"
        ssl_keystore_location => "/opt/logstash/config/server.keystore.p12"
        ssl_keystore_password => "passwd"
        ssl_truststore_location => "/opt/logstash/config/truststore.jks"
        ssl_truststore_password => "passwd"
        #ssl_endpoint_identification_algorithm => ""
    }
}
output {
    elasticsearch {
        hosts => ["172.29.39.141:9200", "172.29.39.142:9200", "172.29.39.143:9200", "172.29.39.144:9200"]
        ssl => true
        cacert => "/etc/certificates/salt-ca.crt"
        user => logstash_internal
        password => "passwd"
        index => "test-1"
    }
}

test2.conf

input {
    kafka {
        bootstrap_servers => "172.29.39.115:9093,172.29.39.116:9093"
        topics => "test-2"
        #consumer_threads => 2
        client_id => "logstash-test-2"
        group_id => "logstash-test-2-0"
        security_protocol => "SSL"
        ssl_keystore_location => "/opt/logstash/config/server.keystore.p12"
        ssl_keystore_password => "passwd"
        ssl_truststore_location => "/opt/logstash/config/truststore.jks"
        ssl_truststore_password => "passwd"
        #ssl_endpoint_identification_algorithm => ""
    }
}
output {
    elasticsearch {
        hosts => ["172.29.39.141:9200", "172.29.39.142:9200", "172.29.39.143:9200", "172.29.39.144:9200"]
        ssl => true
        cacert => "/etc/certificates/salt-ca.crt"
        user => logstash_internal
        password => "passwd"
        index => "test-2"
    }
}

Logstash启动参数

ExecStart=/opt/logstash/bin/logstash "--path.settings" "/opt/logstash/config" "--path.config" "/opt/logstash/config/conf.d/"

问题根因

默认Logstash使用--path.config指定配置目录时,会将该路径下所有配置文件合并为单个管道运行,所有输入源的数据会被分发到该管道下的所有输出节点,因此两个topic的全量数据会同时写入test-1和test-2两个ES索引,导致数据量超出预期。

解决方案
  1. 删除Logstash服务启动参数中的"--path.config" "/opt/logstash/config/conf.d/"
  2. 修改pipelines.yml为每个数据流单独配置独立管道,配置示例如下:
#- pipeline.id: main
  #  path.config: "/opt/logstash/config/conf.d/*.conf"
- pipeline.id: test-1
  path.config: "/opt/logstash/config/conf.d/test1.conf"
  pipeline.workers: 3
- pipeline.id: test-2
  path.config: "/opt/logstash/config/conf.d/test2.conf"
  pipeline.workers: 3

通过多管道配置可以实现不同数据流的完全隔离,避免数据交叉写入错误索引。

内容的提问来源于stack exchange,提问作者mamol

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 08:15:04