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

如何配置Fluentd留存事件字段值并追加至其他源后续事件?

配置Fluentd留存字段值并追加至其他事件

可以实现这个需求,核心思路是通过Lua Filter插件维护全局状态,存储最新的SerialNumber值,在处理不同类型日志时动态更新或追加该字段。

实现步骤

  1. 安装Lua插件
    首先确保Fluentd安装了Lua处理插件:

    gem install fluent-plugin-lua
    
  2. Fluentd配置文件示例
    创建或修改Fluentd配置文件(比如/etc/fluentd/fluent.conf):

    # 输入源:读取目标日志文件,解析日志结构
    <source>
      @type tail
      path /path/to/target/logs.log
      pos_file /var/log/fluentd/log_pos.pos
      tag raw.logs
      <parse>
        @type regexp
        expression /^(?<time>[^ ]+ [^ ]+ [^ ]+) (?<log_type>Log\.[AB]): (?<record>.*)$/
        time_format "%Y-%m-%d %H:%M:%S.%N %z"
      </parse>
    </source>
    
    # 过滤规则:用Lua维护SerialNumber状态并修改日志
    <filter raw.logs>
      @type lua
      script /etc/fluentd/serial_sync.lua
      function filter(tag, timestamp, record)
        -- 初始化全局变量存储最新SerialNumber
        if not _G.latest_serial then
          _G.latest_serial = ""
        end
    
        -- 处理Log.A日志:更新全局SerialNumber
        if record.log_type == "Log.A" then
          local json_data = require("cjson").decode(record.record)
          if json_data.SerialNumber then
            _G.latest_serial = json_data.SerialNumber
          end
          return 1, timestamp, record
        end
    
        -- 处理Log.B日志:追加最新SerialNumber
        if record.log_type == "Log.B" then
          local json_data = require("cjson").decode(record.record)
          json_data.SerialNumber = _G.latest_serial
          record.record = require("cjson").encode(json_data)
          return 1, timestamp, record
        end
    
        -- 其他日志类型直接返回
        return 1, timestamp, record
      end
    </filter>
    
    # 输出目标:根据实际需求替换(比如Elasticsearch、S3等)
    <match raw.logs>
      @type stdout
    </match>
    
  3. Lua脚本文件(可选)
    上面配置中引用的serial_sync.lua内容可单独提取保存,也可直接内嵌在配置中,单独写文件仅为了结构清晰:

    function filter(tag, timestamp, record)
      if not _G.latest_serial then
        _G.latest_serial = ""
      end
    
      if record.log_type == "Log.A" then
        local json_data = require("cjson").decode(record.record)
        if json_data.SerialNumber then
          _G.latest_serial = json_data.SerialNumber
        end
        return 1, timestamp, record
      end
    
      if record.log_type == "Log.B" then
        local json_data = require("cjson").decode(record.record)
        json_data.SerialNumber = _G.latest_serial
        record.record = require("cjson").encode(json_data)
        return 1, timestamp, record
      end
    
      return 1, timestamp, record
    end
    

关键说明

  • 状态维护:利用Lua的全局变量_G.latest_serial存储最新的SerialNumber,每次收到Log.A日志时自动更新。
  • 日志解析:通过正则解析将日志拆分为时间、日志类型和原始JSON记录,方便后续处理。
  • 多进程场景注意:如果Fluentd配置了多worker进程,全局变量会在每个进程中独立存在,此时需要改用外部存储(比如Redis)共享状态,可结合fluent-plugin-redis插件在Lua中读写Redis实现跨进程同步。
  • 初始值处理:首次启动时若未收到Log.A日志,Log.B的SerialNumber会是空字符串,可根据业务需求在Lua脚本中设置默认值。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 11:12:14