如何在Telegraf中转MQTT到InfluxDB时修改行协议并关联MySQL数据
实现方法:在Telegraf中关联MySQL设备元数据并修改InfluxDB行协议
核心思路
利用Telegraf的starlark处理器(支持类Python自定义脚本),在接收MQTT消息后,根据deviceId查询MySQL中的location和area,并将这两个字段添加为InfluxDB行协议的Tag,再转发至InfluxDB。
步骤1:创建Telegraf专用MySQL查询用户
为安全起见,给Telegraf创建仅具备设备表查询权限的MySQL用户:
CREATE USER 'telegraf'@'localhost' IDENTIFIED BY 'your_secure_password'; GRANT SELECT ON your_database.device_table TO 'telegraf'@'localhost'; FLUSH PRIVILEGES;
替换your_secure_password、your_database、device_table为你的实际数据库信息。
步骤2:配置Telegraf的MQTT输入与Starlark处理器
修改Telegraf配置文件(默认路径/etc/telegraf/telegraf.conf),添加/调整以下配置段:
1. MQTT输入插件配置
确保正确配置MQTT输入,能接收原始InfluxDB行协议格式的消息:
[[inputs.mqtt_consumer]] servers = ["tcp://localhost:1883"] topics = ["your/iot/topic"] # 替换为你的实际MQTT主题 data_format = "influx"
2. Starlark处理器配置
添加starlark处理器实现MySQL查询与Tag追加:
[[processors.starlark]] source = ''' from sqlalchemy import create_engine, text # 初始化MySQL连接(仅脚本加载时执行一次) engine = create_engine('mysql+pymysql://telegraf:your_secure_password@localhost/your_database') def apply(metric): # 从当前metric中提取deviceId标签 device_id = metric.tags.get('deviceId') if not device_id: return metric # 查询MySQL获取对应设备的location和area with engine.connect() as conn: query = text("SELECT Location, Area FROM device_table WHERE deviceId = :device_id") result = conn.execute(query, device_id=device_id).fetchone() if result: location, area = result # 将location和area添加为metric的标签 metric.tags['location'] = str(location) metric.tags['area'] = str(area) return metric '''
执行以下命令安装依赖库:
sudo pip3 install pymysql sqlalchemy
3. InfluxDB输出插件配置
确保输出到InfluxDB的配置正确:
[[outputs.influxdb]] urls = ["http://localhost:8086"] database = "your_influx_db" # 替换为你的InfluxDB数据库名 username = "your_influx_user" # 若开启认证请填写 password = "your_influx_password" # 若开启认证请填写
步骤3:测试与重启Telegraf
- 验证配置文件合法性:
telegraf --config /etc/telegraf/telegraf.conf --test
- 重启Telegraf服务:
sudo systemctl restart telegraf
替代方案:静态设备元数据的无脚本实现
如果设备的location和area不会频繁变更,可采用以下无脚本方案:
- 用
mysql输入插件定期拉取设备元数据到InfluxDB:
[[inputs.mysql]] servers = ["telegraf:your_secure_password@tcp(localhost:3306)/your_database"] queries = [ "SELECT deviceId, Location as location, Area as area FROM device_table" ] measurement_name = "device_metadata"
- 使用
join处理器关联MQTT数据与设备元数据:
[[processors.join]] [[processors.join.left]] measurement = "your_original_measurement" # MQTT输入对应的measurement名称 tags = ["deviceId"] [[processors.join.right]] measurement = "device_metadata" tags = ["deviceId"] [[processors.join.as]] measurement = "your_original_measurement"
结果验证
发送测试MQTT消息:
mosquitto_pub -t "your/iot/topic" -m "measurement,deviceId=\"A\" temperature=\"25\",humidity=\"60\""
在InfluxDB中执行查询:
SELECT * FROM measurement WHERE deviceId='A'
可看到location和area标签已成功追加到数据中。
内容的提问来源于stack exchange,提问作者Gimhana Jayasekara
相关产品推荐
相关产品推荐

