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

Logstash MQTT输出断开后无法自动重连,需重启容器求助

问题:Logstash MQTT连接中断后管道终止,需手动重启容器恢复

部署的Logstash通过MQTT传输API数据,当MQTT连接中断时,数据管道直接终止运行,必须重启Logstash容器才能恢复。已尝试添加健康检查但未生效,需要实现连接异常时自动触发管道重启。

现有配置

docker-compose.yml

logstash:
  build: ./logstash
  image: docker.elastic.co/logstash/logstash/wins:8.8.2
  container_name: logstash_hktm_multi
  volumes:
    - type: bind
      source: ./logstash/config/logstash.yml
      target: /usr/share/logstash/config/logstash.yml
      read_only: true
    - type: bind
      source: ./logstash/config/pipelines.yml
      target: /usr/share/logstash/config/pipelines.yml
      read_only: true
    - type: bind
      source: ./logstash/pipeline
      target: /usr/share/logstash/pipeline
      read_only: true
    - type: bind
      source: ./logstash/data
      target: /var/lib/logstash/data
  ports:
    - 5044:5044
    - 5000:5000/tcp
    - 5000:5000/udp
    - 9600:9600
  command: --config.reload.automatic 
  environment:
    LS_JAVA_OPTS: -Xmx1g -Xms1g
    LOGSTASH_INTERNAL_PASSWORD: ${LOGSTASH_INTERNAL_PASSWORD:-}
    TZ: Europe/Istanbul
    start_date: '1900-01-01'
  healthcheck:
    test: ["CMD", "curl", "-f", "http://localhost:9600"]
    interval: 3m
    timeout: 60s
    retries: 5
          
  networks:
    - elk
  depends_on:
    - elasticsearch
  restart: always

Logstash管道配置

input {
  http_poller {
    urls => {
      dummy_request => {
        method => get
        url => "http://localhost:9600" 
      }
    }
    request_timeout => 60
    codec => "json"
    schedule => { "every" => "1m" }
  }
}

filter {
  ruby {
    code => 'event.set("date", Time.now.strftime("%Y-%m-%d"))'
  }

  ruby {
    code => 'logger.info("Sistem saati: #{Time.now.strftime("%Y-%m-%d %H:%M:%S")}")'
  }

  http {
    url => "https:/xxx"
    target_body => "response"
  }

  ruby {
    code => "event.set('total_energy', event.get('[response][energy][values][0][value]') || 0)"
  }

  mutate {  
    remove_field => [ "@version","event", "id", "name","ephemeral_id","status","snapshot", "response", "version","http_address","pipeline","monitoring","build_date","build_sha","build_snapshot", "host", "url", "error", "http", "tags", "dummy_request" ]
    add_field =>  { "isIncrement" => "true" }
    add_field =>  { "timeUnit" => "DAY" }
  }
}

output {
  mqtt {
    host => "192.168.2.205"
    port => "1883"
    topic => "ry_GESS"
  }
  stdout {
    codec => rubydebug
  }
}

连接异常报错日志

[2024-06-25T17:37:10,276][INFO ][logstash.filters.ruby    ][solaredge][f936179747d4760133b949bcf0600c72e081a7c5cb61ca44969561336c404f8b] Sistem saati: 2024-06-25 17:37:10
[2024-06-25T17:37:10,522]

[ERROR][logstash.javapipeline    ][solaredge] Pipeline worker error, the pipeline will be stopped {:pipeline_id=>"solaredge", :error=>"

(NotConnectedException) MQTT::NotConnectedException", :exception=>Java::OrgJrubyExceptions::Exception, :backtrace=>["RUBY.send_packet(/usr/share/logstash/vendor/bundle/jruby/3.1.0/gems/mqtt-0.6.0/lib/mqtt/client.rb:567)", 


"RUBY.publish(/usr/share/logstash/vendor/bundle/jruby/3.1.0/gems/mqtt-0.6.0/lib/mqtt/client.rb:330)", "RUBY.handle_events(/usr/share/logstash/vendor/bundle/jruby/3.1.0/gems/logstash-output-mqtt-1.2.1/lib/logstash/outputs/mqtt.rb:176)", "RUBY.multi_receive(/usr/share/logstash/vendor/bundle/jruby/3.1.0/gems/logstash-output-mqtt-1.2.1/lib/logstash/outputs/mqtt.rb:158)", 


"org.logstash.config.ir.compiler.AbstractOutputDelegatorExt.multi_receive(org/logstash/config/ir/compiler/AbstractOutputDelegatorExt.java:121)", "RUBY.start_workers(/usr/share/logstash/logstash-core/lib/logstash/java_pipeline.rb:304)"], :thread=>"#<Thread:0x30c8f4a8 /usr/share/logstash/logstash-core/lib/logstash/java_pipeline.rb:134 sleep>"}

解决方案

1. 优化MQTT输出插件的重连与容错配置

当前MQTT输出默认未开启自动重连,连接中断后直接导致管道崩溃。修改管道配置中的MQTT输出,添加重连相关参数:

output {
  mqtt {
    host => "192.168.2.205"
    port => "1883"
    topic => "ry_GESS"
    # 自动重连间隔(秒)
    reconnect_interval => 5
    # 无限重试重连(0表示无限制)
    reconnect_limit => 0
    # 发送超时时间,避免阻塞管道
    send_timeout => 10
    # 设置固定客户端ID,配合持久会话
    client_id => "logstash_solaredge"
    # 启用持久会话,重连后恢复状态
    clean_session => false
  }
  stdout {
    codec => rubydebug
  }
}

2. 改进Docker健康检查,精准检测管道状态

原健康检查仅检测Logstash管理端口,无法感知管道是否正常运行。修改docker-compose.yml中的健康检查规则:

healthcheck:
  test: ["CMD", "curl", "-f", "http://localhost:9600/_node/pipelines/solaredge"]
  interval: 1m
  timeout: 10s
  retries: 3
  start_period: 30s

注意:solaredge需与pipelines.yml中定义的管道ID一致。若管道停止,该接口会返回非200状态码,触发健康检查失败。

3. 启用Logstash管道自动重启功能

在logstash.yml中添加以下配置,让管道崩溃后自动重启:

pipeline:
  # 开启管道自动重启
  auto_restart: true
  # 最大重启尝试次数
  restart_retry_count: 5
  # 重启间隔(秒)
  restart_backoff_interval: 5

4. 容器级兜底重启策略

保留docker-compose.yml中的重启配置,改为基于健康检查失败的重启:

restart: on-failure:5

说明:当健康检查连续失败5次时,Docker自动重启容器。

错误分析

从报错日志可知,MQTT连接中断抛出NotConnectedException后,Logstash管道worker直接崩溃,且无自动重启机制。核心问题是MQTT输出插件未配置重连逻辑,同时管道未开启自动重启。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:03:11