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

如何用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': []}"

正确配置方案

配置说明

  1. 替换inputs.exec的解析格式为json_v2,灵活提取嵌套数组内的内容
  2. 使用processors.split拆分field2的数组,将每个元素转为单独的消息
  3. 为两个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 23:27:38