如何用Telegraf单inputs.exec向多outputs.kafka分字段发数据?
问题:Telegraf拆分JSON字段发送到不同Kafka主题
我通过Telegraf的inputs.exec插件采集到包含多字段的JSON数据,格式如下:
{ "field1": [ { "abc" : 0, "efg" : 1, "hij" : 4, "jkl" : 5 } ], "field2": [ { "host" : "admin1", "timestamp": 1682314679774 }, { "host" : "admin2", "timestamp": 1682314679775 }, { "host" : "admin3", "timestamp": 1682314679773 } ] }
预期目标
- 将
field1数组内的JSON对象发送至Kafka的field1主题,输出内容为:{"abc" : 0,"efg" : 1,"hij" : 4,"jkl" : 5} - 将
field2数组内的每个JSON对象分别发送至Kafka的field2主题,每条消息对应一个对象:{"host" : "admin1","timestamp": 1682314679774}{"host" : "admin2","timestamp": 1682314679775}{"host" : "admin3","timestamp": 1682314679773}
尝试的错误配置
[[outputs.kafka]] brokers = ["admin:9091"] topic = "field1" data_format = "json" flush_interval = "1m" [[outputs.kafka]] brokers = ["admin:9092"] ## Kafka topic for producer messages topic = "field2" data_format = "json" flush_interval = "1m" [[inputs.exec]] interval = "1m" commands = ["/opt/clustertest/bin/script.py"] timeout = "10s" # I want the script to execute every ten seconds data_format = "json" flush_interval = "1m" json_query = "{'field1': [], 'field2': []}"
正确配置方案
配置说明
- 替换
inputs.exec的解析格式为json_v2,灵活提取嵌套数组内的内容 - 使用
processors.split拆分field2的数组,将每个元素转为单独的消息 - 为两个Kafka输出配置字段过滤,确保仅发送对应主题需要的数据
完整配置
# 输入插件:执行脚本采集JSON数据,10秒执行一次 [[inputs.exec]] interval = "10s" # 匹配你想要的10秒执行间隔 commands = ["/opt/clustertest/bin/script.py"] timeout = "10s" data_format = "json_v2" # 解析JSON,提取field1和field2的内容 [[inputs.exec.json_v2]] # 提取field1数组的第一个元素,重命名为field1_data [[inputs.exec.json_v2.field]] path = "field1[0]" rename = "field1_data" # 提取field2整个数组,重命名为field2_data [[inputs.exec.json_v2.field]] path = "field2" rename = "field2_data" # 处理器:拆分field2_data数组,每个元素生成一条独立消息 [[processors.split]] [[processors.split.field]] name = "field2_data" tag = "field2_index" # 可选,标记拆分后的元素索引,可删除 # 输出插件1:发送field1数据到field1主题 [[outputs.kafka]] brokers = ["admin:9091"] topic = "field1" data_format = "json" flush_interval = "1m" # 仅序列化field1_data字段,输出纯净的JSON对象 json_fields = ["field1_data"] # 过滤仅保留包含field1_data的消息 [[outputs.kafka.tagpass]] field1_data = [".+"] # 输出插件2:发送field2拆分后的数据到field2主题 [[outputs.kafka]] brokers = ["admin:9092"] topic = "field2" data_format = "json" flush_interval = "1m" # 仅序列化field2_data字段,输出纯净的JSON对象 json_fields = ["field2_data"] # 过滤仅保留包含field2_data的消息 [[outputs.kafka.tagpass]] field2_data = [".+"]
错误点修正
- 移除了
inputs.exec中的flush_interval:该参数属于输出插件,输入插件不支持 - 修正了执行间隔:将
interval改为10s,匹配你注释中想要的10秒执行频率 - 替换
data_format为json_v2:原生json格式无法灵活处理嵌套数组的提取,json_v2支持JSONPath语法精准提取字段 - 添加
processors.split:解决field2数组无法拆分为单条消息的问题 - 为Kafka输出添加字段过滤:避免所有数据被发送到两个主题,确保数据流向正确
内容的提问来源于stack exchange,提问作者Groot
相关产品推荐
相关产品推荐

