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

基于Kafka驱动的Elasticsearch文档迁移(索引A→B)的Logstash配置报错排查

基于Kafka驱动的Elasticsearch文档迁移(索引A→B)的Logstash配置报错排查

问题描述

我正在尝试用Logstash实现一套Elasticsearch文档迁移流程,具体逻辑是:

  • 从Kafka消费包含目标文档ID的消息
  • 用这个ID去Elasticsearch索引A中查询对应文档
  • 将查询到的文档写入Elasticsearch索引B
  • 最后删除索引A中的原文档

但配置完成后启动Logstash时遇到了连接错误,想请大家帮忙排查问题。

当前使用的Logstash配置

input {
  kafka {
    bootstrap_servers => "ip:port"
    topics => "topic"
    group_id => "topic-19092024"
    consumer_threads => 2
    max_poll_interval_ms => "300000"
    ssl_keystore_location => "path/kafka.keystore.jks"
    ssl_keystore_password => "abc"
    ssl_truststore_location => "path/kafka.truststore.jks"
    ssl_truststore_password => "abc"
    ssl_endpoint_identification_algorithm => ""
    security_protocol => "SSL"
  }
}

filter {
  # Parse Kafka JSON message to extract 'id'
  json {
    source => "message"
    target => "eventdata"
  }

  # Get document by _id from index A
  elasticsearch {
    hosts => ["https://ip:9200"]
    index => "index-a"
    user => "${ES_USER}"
    password => "${ES_PWD}"
    ca_file => "path/elastic-certificate.pem"
    query => "_id:%{[eventdata][id]}"
    fields => { "[_source]" => "esdoc" }
    ssl => true
  }

  # Remove events where document does not exist on Elasticsearch
  if ![esdoc] {
    drop { }
  } else {
    # Move all fields of [esdoc] to root (flatten)
    ruby {
      code => '
        if event.get("esdoc")
          event.get("esdoc").each {|k,v| event.set(k, v)}
        end
      '
    }
  }
}

output {
  # 1. Write document to index B (with same id)
  elasticsearch {
    hosts => ["https://ip:9200"]
    index => "index-b"
    document_id => "%{[eventdata][id]}"
    user => "${ES_USER}"
    password => "${ES_PWD}"
    cacert => "path/elastic-certificate.pem"
    ssl => true
    ssl_certificate_verification => false
    action => "index"
  }

  # 2. Delete document from index A
  elasticsearch {
    hosts => ["https://ip:9200"]
    index => "index-a"
    document_id => "%{[eventdata][id]}"
    user => "${ES_USER}"
    password => "${ES_PWD}"
    cacert => "path/elastic-certificate.pem"
    ssl => true
    ssl_certificate_verification => false
    action => "delete"
  }
}

报错信息

:exception=>#<Manticore::ResolutionFailure: https>, :backtrace=>[
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/manticore-0.7.1-java/lib/manticore/response.rb:36:in block in initialize'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/manticore-0.7.1-java/lib/manticore/response.rb:79:in call'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/manticore-0.7.1-java/lib/manticore/response.rb:274:in call_once'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/manticore-0.7.1-java/lib/manticore/response.rb:158:in code'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/elasticsearch-transport-7.15.0/lib/elasticsearch/transport/transport/http/manticore.rb:103:in block in perform_request'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/elasticsearch-transport-7.15.0/lib/elasticsearch/transport/transport/base.rb:286:in perform_request'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/elasticsearch-transport-7.15.0/lib/elasticsearch/transport/transport/http/manticore.rb:85:in perform_request'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/elasticsearch-transport-7.15.0/lib/elasticsearch/transport/client.rb:191:in perform_request'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/elasticsearch-7.15.0/lib/elasticsearch.rb:92:in elasticsearch_validation_request'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/elasticsearch-7.15.0/lib/elasticsearch.rb:51:in verify_elasticsearch'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/elasticsearch-7.15.0/lib/elasticsearch.rb:40:in method_missing'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/elasticsearch-api-7.15.0/lib/elasticsearch/api/actions/ping.rb:38:in ping'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/logstash-filter-elasticsearch-3.11.0/lib/logstash/filters/elasticsearch.rb:324:in test_connection!'",
"/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/logstash-filter-elasticsearch-3.11.0/lib/logstash/filters/elasticsearch.rb:113:in register'",
"org/logstash/config/ir/compiler/AbstractFilterDelegatorExt.java:75:in register'",
"/usr/share/logstash/logstash-core/lib/logstash/java_pipeline.rb:232:in block in register_plugins'",
"org/jruby/RubyArray.java:1821:in each'",
"/usr/share/logstash/logstash-core/lib/logstash/java_pipeline.rb:231:in register_plugins'",
"/usr/share/logstash/logstash-core/lib/logstash/java_pipeline.rb:590:in maybe_setup_out_plugins'"
]

排查思路与修复建议

这个报错是Manticore::ResolutionFailure,本质是Logstash无法解析或连接到Elasticsearch地址,建议从以下几个方向排查:

1. 验证Elasticsearch地址的可达性

  • 首先确认配置中hosts字段的IP和端口是否正确,没有拼写错误
  • 在Logstash所在服务器上,用命令测试ES连接:
    curl -u ${ES_USER}:${ES_PWD} https://ip:9200 --cacert path/elastic-certificate.pem
    
    如果连接失败,检查网络防火墙/安全组是否开放了9200端口,ES节点是否正常运行

2. 统一SSL配置参数,避免不一致

  • 注意到你在filter的elasticsearch插件中用了ca_file,而output中用了cacert,这两个参数是等效的,但建议统一写法,避免混淆
  • 在filter的elasticsearch配置中添加ssl_certificate_verification => false,和output保持一致,先排除证书验证的问题:
    elasticsearch {
      hosts => ["https://ip:9200"]
      index => "index-a"
      user => "${ES_USER}"
      password => "${ES_PWD}"
      ca_file => "path/elastic-certificate.pem"
      query => "_id:%{[eventdata][id]}"
      fields => { "[_source]" => "esdoc" }
      ssl => true
      ssl_certificate_verification => false # 添加这一行
    }
    

3. 确认环境变量加载正常

  • 配置中使用了${ES_USER}和${ES_PWD}环境变量,要确保Logstash启动时这些变量已经正确设置
  • 可以在启动Logstash时添加--verbose参数,查看日志中是否有变量加载的相关信息,或者直接替换成硬编码的账号密码做测试(测试完成后记得改回环境变量)

4. 调整连接超时参数

  • 如果是网络延迟导致的连接失败,可以在elasticsearch配置中添加超时参数:
    timeout => 30 # 单位:秒
    

修复后的参考配置

input {
  kafka {
    bootstrap_servers => "ip:port"
    topics => "topic"
    group_id => "topic-19092024"
    consumer_threads => 2
    max_poll_interval_ms => "300000"
    ssl_keystore_location => "path/kafka.keystore.jks"
    ssl_keystore_password => "abc"
    ssl_truststore_location => "path/kafka.truststore.jks"
    ssl_truststore_password => "abc"
    ssl_endpoint_identification_algorithm => ""
    security_protocol => "SSL"
  }
}

filter {
  json {
    source => "message"
    target => "eventdata"
  }

  elasticsearch {
    hosts => ["https://ip:9200"]
    index => "index-a"
    user => "${ES_USER}"
    password => "${ES_PWD}"
    ca_file => "path/elastic-certificate.pem"
    query => "_id:%{[eventdata][id]}"
    fields => { "[_source]" => "esdoc" }
    ssl => true
    ssl_certificate_verification => false
    timeout => 30
  }

  if ![esdoc] {
    drop { }
  } else {
    ruby {
      code => '
        if event.get("esdoc")
          event.get("esdoc").each {|k,v| event.set(k, v)}
        end
      '
    }
  }
}

output {
  elasticsearch {
    hosts => ["https://ip:9200"]
    index => "index-b"
    document_id => "%{[eventdata][id]}"
    user => "${ES_USER}"
    password => "${ES_PWD}"
    ca_file => "path/elastic-certificate.pem" # 统一为ca_file
    ssl => true
    ssl_certificate_verification => false
    action => "index"
    timeout => 30
  }

  elasticsearch {
    hosts => ["https://ip:9200"]
    index => "index-a"
    document_id => "%{[eventdata][id]}"
    user => "${ES_USER}"
    password => "${ES_PWD}"
    ca_file => "path/elastic-certificate.pem" # 统一为ca_file
    ssl => true
    ssl_certificate_verification => false
    action => "delete"
    timeout => 30
  }
}

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 10:09:30