基于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连接:
如果连接失败,检查网络防火墙/安全组是否开放了9200端口,ES节点是否正常运行curl -u ${ES_USER}:${ES_PWD} https://ip:9200 --cacert path/elastic-certificate.pem
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
相关产品推荐
相关产品推荐

