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

FluentD in_forward插件误将JSON对象识别为数组并报'Incoming chunk is broken'错误的排查求助

FluentD in_forward插件误将JSON对象识别为数组并报'Incoming chunk is broken'错误的排查求助

大家好,我最近在AWS环境的FluentD节点上遇到了一个棘手的问题,折腾了好几天没找到根源,想请教下各位大佬:

我用Ruby代码通过TLS连接向FluentD的in_forward插件发送格式完好的JSON对象,TLS连接是成功的,但FluentD一直弹出**"Incoming chunk is broken"**的警告。更奇怪的是,明明发送的是单个JSON对象,代码里msg.is_a?(Array)居然返回true,而且FluentD日志里同时显示了原始JSON字符串和转换成Ruby哈希的内容(冒号被替换成了=>)。

环境信息

  • Ruby版本:2.7.0p0 (2019-12-25 revision 647ee6f091) [x86_64-linux-gnu]
  • FluentD使用插件:in_forward

发送端Ruby代码

begin
  socket = TCPSocket.new(host, port)

  ssl_context = OpenSSL::SSL::SSLContext.new()
  rsa_cert = OpenSSL::X509::Certificate.new(ssl_cert_string)
  rsa_pkey = OpenSSL::PKey::RSA.new(ssl_key_string)
  ca_intermediate_cert = OpenSSL::X509::Certificate.new(ssl_ca_cert_string)

  ssl_context.add_certificate(rsa_cert, rsa_pkey, [ca_intermediate_cert])
  ssl_context.verify_mode = OpenSSL::SSL::VERIFY_NONE
  ssl_context.ssl_version = :TLSv1_2

  ssl_socket = OpenSSL::SSL::SSLSocket.new(socket, ssl_context)
  ssl_socket.sync_close = true
  ssl_socket.connect

  ssl_socket.puts(msg)
rescue Exception => e
  puts "FluentD exception :#{e.message}"
ensure
  ssl_socket&.close
  socket&.close
end

FluentD配置(fluent.conf相关片段)

<source>
  @type forward
  port {{port}}
  bind 0.0.0.0
  <parse>
    @type json
  </parse>
  <transport tls>
    version TLSv1_2
    cert_path {{tls_certificate}}
    private_key_path {{tls_private_key}}
    private_key_passphrase ""
  </transport>
  tag fluentd_ssl
</source>

并发调用逻辑

foo = ProgressBar.new(N*M)
start_concurrent = Time.now

N.times.map do
  foo.increment!
  bar = ProgressBar.new(M)
  Thread.new do
    M.times do
      bar.increment!
      send_message(host, port, ssl_ca_cert_string, ssl_cert_string, ssl_key_string)
      sleep 0.1
    end
  end
end.each(&:join)

finish_concurrent = Time.now
duration = humanize(finish_concurrent - start_concurrent)
puts "Concurrent duration: #{duration}"

已尝试的解决方法

  • 替换Socket写入方法:原本用ssl_socket.puts(msg),后来换成ssl_socket.syswrite(msg.to_s)和ssl_socket.write(msg.to_s),但问题依然存在
  • 排查消息拆分问题:怀疑是chunk size小于消息大小导致被拆分成数组,但FluentD日志里能看到完整的JSON对象,所以应该不是这个原因

核心疑问

  • 为什么明明发送的是单个JSON对象,msg.is_a?(Array)会返回true?
  • 不使用那个猴子补丁的话,有没有其他方法可以排查"Incoming chunk is broken"的错误?日志里提到是msgpack解包器收到了损坏的流,但我发送的是JSON啊
  • 如果必须用那个猴子补丁来添加自定义日志(想在FluentD里打印on_message的内容),作为Ruby新手,我该怎么把它应用到FluentD上?

麻烦各位帮忙分析下可能的原因,或者给我一些排查方向,谢谢!

备注:内容来源于stack exchange,提问作者neuralsea

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.23 07:14:15