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

Logstash单Kafka输入按项目分流多输出(文件/数据库)方案咨询

需求可行性结论

该需求完全可以实现,核心是利用Logstash的Filter阶段做字段判断打标签,再在Output阶段基于标签做条件分流即可。

具体实现思路

首先确认Kafka消息中用来区分项目维度的业务字段,比如假设data.projectId字段值为1代表项目1,值为2代表项目2,按如下步骤修改配置即可:

  1. 新增Filter配置段,给不同项目的数据打识别标签
filter {
  # 识别项目1数据,添加专属标签
  if [data][projectId] == 1 {
    mutate {
      add_tag => ["project1"]
    }
  }
  # 识别项目2数据,添加专属标签
  else if [data][projectId] == 2 {
    mutate {
      add_tag => ["project2"]
    }
  }
}

注意:上面配置中用来判断项目的[data][projectId]字段请替换为你实际业务中用来区分两个项目的字段,字段值也对应调整为你们两个项目的实际标识值。

  1. 修改Output配置段,按标签做分流输出
output {
  # 全量数据保持写入ES的原有逻辑不变
  elasticsearch {
    hosts => ["127.0.0.1:9200"]
    index => "survey-result"
    workers => 1
    codec => json
  }
  # 仅项目1的数据输出txt和csv文件
  if "project1" in [tags] {
    file {
      path => "/sadata/saexport-%{+YYYY-MM-dd}.txt"
    }
    csv {
      fields => ["data.integrationName", "data.contextData.UCID", "data.contextData.ANI", "data.reason.reason", "data.status"]
      path => "/sadata/saexport.%{+YYYY-MM-dd-hh}.csv"
    }
  }
  # 仅项目2的数据写入指定数据库,这里以MySQL为例使用jdbc输出插件
  if "project2" in [tags] {
    jdbc {
      driver_class => "com.mysql.cj.jdbc.Driver"
      connection_string => "jdbc:mysql://数据库地址:3306/目标库名?useSSL=false&serverTimezone=Asia/Shanghai"
      username => "数据库账号"
      password => "数据库密码"
      statement => "INSERT INTO 目标表名 (字段1,字段2) VALUES (?,?)"
      # 按顺序对应SQL中问号的取值,替换为实际业务字段即可
      parameters => [ "字段1", "%{[data][xxx]}", "字段2", "%{[data][yyy]}" ]
    }
  }
}

注意:使用jdbc输出插件前需要提前将对应数据库的JDBC驱动jar包放到Logstash的logstash-core/lib/jars/目录下,重启Logstash后生效。

配置参考资料获取渠道

  • Logstash安装包自带的本地文档,部署后可直接在安装目录的docs文件夹下访问,包含所有内置插件的完整参数说明
  • Logstash配置目录下的官方示例文件,包含filter、各类input/output插件的真实使用样例
  • 企业内部已上线的同类型Logstash分流配置案例,可直接复用适配

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.05 16:21:00