如何实现InfluxDB数据向MQTT Broker/服务器的传输
问题:实现InfluxDB到MQTT的数据传输桥接
这是我首次开展MQTT相关开发工作,最终目标是将InfluxDB中的数据同步至Snowflake,在实现该目标前需要先完成如下前置任务:
- 实现InfluxDB到MQTT的数据传输,目前未检索到公开的相关实现示例。
我此前已经实现过MQTT数据存储到InfluxDB的功能,使用的脚本如下:
"""A MQTT to InfluxDB Bridge This script receives MQTT data and saves those to InfluxDB. """ import re from typing import NamedTuple import paho.mqtt.client as mqtt from influxdb import InfluxDBClient INFLUXDB_ADDRESS = '10.10.10.247' INFLUXDB_USER = 'iotuser' INFLUXDB_PASSWORD = 'iotpassword' INFLUXDB_DATABASE = 'homeiot_db' MQTT_ADDRESS = '10.10.10.247' MQTT_USER = 'iotuser' MQTT_PASSWORD = 'iotpassword' MQTT_TOPIC = 'home/+/+' # [room]/[temperature|humidity|light|status] MQTT_REGEX = 'home/([^/]+)/([^/]+)' MQTT_CLIENT_ID = 'MQTTInfluxDBBridge' influxdb_client = InfluxDBClient(INFLUXDB_ADDRESS, 8086, INFLUXDB_USER, INFLUXDB_PASSWORD, None) class SensorData(NamedTuple): location: str measurement: str value: float def on_connect(client, userdata, flags, rc): """ The callback for when the client receives a CONNACK response from the server.""" print('Connected with result code ' + str(rc)) client.subscribe(MQTT_TOPIC) def on_message(client, userdata, msg): """The callback for when a PUBLISH message is received from the server.""" print(msg.topic + ' ' + str(msg.payload)) sensor_data = _parse_mqtt_message(msg.topic, msg.payload.decode('utf-8')) if sensor_data is not None: _send_sensor_data_to_influxdb(sensor_data) def _parse_mqtt_message(topic, payload): match = re.match(MQTT_REGEX, topic) if match: location = match.group(1) measurement = match.group(2) if measurement == 'status': return None return SensorData(location, measurement, float(payload)) else: return None def _send_sensor_data_to_influxdb(sensor_data): json_body = [ { 'measurement': sensor_data.measurement, 'tags': { 'location': sensor_data.location }, 'fields': { 'value': sensor_data.value } } ] print (json_body) influxdb_client.write_points(json_body) def _init_influxdb_database(): databases = influxdb_client.get_list_database() if len(list(filter(lambda x: x['name'] == INFLUXDB_DATABASE, databases))) == 0: influxdb_client.create_database(INFLUXDB_DATABASE) influxdb_client.switch_database(INFLUXDB_DATABASE) def main(): _init_influxdb_database() mqtt_client = mqtt.Client(MQTT_CLIENT_ID) mqtt_client.username_pw_set(MQTT_USER, MQTT_PASSWORD) mqtt_client.on_connect = on_connect mqtt_client.on_message = on_message mqtt_client.connect(MQTT_ADDRESS, 1883) mqtt_client.loop_forever() if __name__ == '__main__': print('MQTT to InfluxDB bridge') main()
如果有开发者曾实现过该需求,或知晓可行的实现方案,希望能得到相关的帮助与指导。
回答
你现有脚本是单向从MQTT消费消息写入InfluxDB,反向实现InfluxDB到MQTT的传输不需要找现成的完整示例,核心逻辑拆成两部分即可:从InfluxDB获取增量数据,调用MQTT客户端发布消息,复用你脚本里已经写好的连接配置就能快速跑通。
增量数据获取方案二选一
- 生产低延迟场景:用InfluxDB 1.x自带的Subscriber功能,在InfluxDB配置文件中添加订阅规则,所有新写入的数据点会主动推送到你指定的HTTP接口,不需要轮询查库,延迟可做到毫秒级。
- 快速验证/小规模场景:用时间戳标记位点轮询查询,每次拉取大于上次同步时间的新数据,拉取完成后更新同步位点,不需要修改InfluxDB任何配置,实现成本最低。
参考实现代码
下面的代码基于你现有的配置编写,用轮询方案实现,直接合并到原有脚本中替换main函数即可使用:
import time from datetime import datetime # 新增配置 SYNC_POS_FILE = './last_sync_pos.txt' QUERY_INTERVAL = 5 # 增量查询间隔,单位秒 INFLUXDB_MQTT_CLIENT_ID = 'InfluxDBMQTTBridge' def get_last_sync_pos(): """读取上次同步的时间位点,首次运行默认取当前时间前1小时""" try: with open(SYNC_POS_FILE, 'r') as f: return datetime.fromisoformat(f.read().strip()) except FileNotFoundError: return datetime.fromtimestamp(time.time() - 3600) def save_sync_pos(pos): """持久化最新同步位点""" with open(SYNC_POS_FILE, 'w') as f: f.write(pos.isoformat()) def query_influx_incremental(start_time): """查询指定时间点之后的所有传感器数据""" query_sql = f""" SELECT * FROM /.*/ WHERE time > '{start_time.isoformat()}' """ result = influxdb_client.query(query_sql, database=INFLUXDB_DATABASE) publish_list = [] latest_time = start_time for measurement_name, point_series in result.items(): for point in point_series: # 按原有Topic规则组装:home/{location}/{measurement} location = point.get('location', 'default_room') value = point['value'] target_topic = f'home/{location}/{measurement_name}' publish_list.append((target_topic, str(value))) # 更新最新时间位点 point_time = datetime.fromisoformat(point['time'].replace('Z', '+00:00')).replace(tzinfo=None) if point_time > latest_time: latest_time = point_time return publish_list, latest_time def sync_worker(mqtt_client): last_sync_pos = get_last_sync_pos() while True: pending_publish, latest_pos = query_influx_incremental(last_sync_pos) for topic, payload in pending_publish: # QoS设为1保证消息不丢 mqtt_client.publish(topic, payload, qos=1) print(f"Published {payload} to {topic}") if latest_pos > last_sync_pos: save_sync_pos(latest_pos) last_sync_pos = latest_pos time.sleep(QUERY_INTERVAL) def main(): _init_influxdb_database() # 初始化MQTT客户端,注意客户端ID不要和原有桥接脚本重复 mqtt_client = mqtt.Client(INFLUXDB_MQTT_CLIENT_ID) mqtt_client.username_pw_set(MQTT_USER, MQTT_PASSWORD) mqtt_client.connect(MQTT_ADDRESS, 1883) # 启动MQTT后台网络循环 mqtt_client.loop_start() # 启动增量同步任务 sync_worker(mqtt_client) if __name__ == '__main__': print('InfluxDB to MQTT bridge started') main()
注意事项
- 新脚本的MQTT客户端ID不要和之前的
MQTTInfluxDBBridge重复,否则两个客户端连接同一个MQTT服务端会互相踢下线。 - 轮询间隔可以根据你的业务对延迟的容忍度调整,最低可以设到1秒,不会对InfluxDB造成明显压力。
- 等InfluxDB到MQTT的链路跑通,后续同步到Snowflake直接消费对应MQTT主题的数据即可,Snowflake原生支持MQTT数据源接入,不需要额外做复杂的格式转换。
- 如果后续要换成Subscriber推送方案,只需要把轮询拉取数据的部分替换成HTTP接口接收InfluxDB推送的点即可,MQTT发布的逻辑完全不用修改。
内容的提问来源于stack exchange,提问作者oussama AY
相关产品推荐
相关产品推荐

