Logstash单Kafka输入按项目分流多输出(文件/数据库)方案咨询
需求可行性结论
该需求完全可以实现,核心是利用Logstash的Filter阶段做字段判断打标签,再在Output阶段基于标签做条件分流即可。
具体实现思路
首先确认Kafka消息中用来区分项目维度的业务字段,比如假设data.projectId字段值为1代表项目1,值为2代表项目2,按如下步骤修改配置即可:
- 新增Filter配置段,给不同项目的数据打识别标签
filter { # 识别项目1数据,添加专属标签 if [data][projectId] == 1 { mutate { add_tag => ["project1"] } } # 识别项目2数据,添加专属标签 else if [data][projectId] == 2 { mutate { add_tag => ["project2"] } } }
注意:上面配置中用来判断项目的
[data][projectId]字段请替换为你实际业务中用来区分两个项目的字段,字段值也对应调整为你们两个项目的实际标识值。
- 修改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
相关产品推荐
相关产品推荐

