使用MQTT发送1秒批次相机遥测数据如何避免数据丢失
MQTT批量上报数据优化方案
原代码核心问题
- 用
dict.update()合并数据:相同键名的遥测数据会被直接覆盖,是数据丢失的核心原因 response.fp.readline()为阻塞调用:窗口时间到期时如果正在阻塞读流,会导致上报延迟,甚至错过1秒上报要求- 重复执行
data.decode('utf-8'):既浪费性能,也可能因二进制数据异常导致重复解码报错 - 无异常捕获与数据兜底逻辑:流读取失败、解码失败、MQTT发布失败都会直接终止程序,未发送的缓存数据直接丢失
- QoS=1发布无回执确认:MQTT QoS1需要等待broker回执确认发布成功,原代码未做确认逻辑,发布失败数据直接丢失
优化实现代码
import time import json import select import os from collections import deque # 本地兜底文件,用于存储发布失败的缓存数据 BACKUP_FILE = "mqtt_backup.json" # 补发队列:启动时优先读取兜底文件的未发送数据 pending_send = deque() if os.path.exists(BACKUP_FILE): try: with open(BACKUP_FILE, "r", encoding="utf-8") as f: pending_send.extend(json.load(f)) os.remove(BACKUP_FILE) except: pass # MQTT发布成功回调,确认数据上报成功 publish_success = False def on_publish(client, userdata, mid): global publish_success publish_success = True myMQTTClient.on_publish = on_publish while True: t_end = time.time() + 1 # 用列表存储1秒内的所有有效数据,避免同键覆盖 payload_list = [] while time.time() < t_end: timeout = t_end - time.time() if timeout <= 0: break # 用select监听流可读状态,超时自动退出避免阻塞 readable, _, _ = select.select([response.fp], [], [], timeout) if not readable: continue data = response.fp.readline() try: data_str = data.decode("utf-8").strip() if data_str.startswith("data"): raw = data_str[5:] # 先转成结构化数据再存入列表 raw_data = json.loads(raw) payload_list.append(raw_data) except (UnicodeDecodeError, json.JSONDecodeError): # 异常数据可自行选择打日志留存,此处直接跳过 continue # 合并待补发数据与当前窗口数据 if pending_send: payload_list = list(pending_send) + payload_list pending_send.clear() if payload_list: publish_success = False # 上报数据增加窗口标识,方便下游排查问题 publish_payload = json.dumps({ "window_timestamp": int(time.time()), "camera_count": 3, "data_list": payload_list }) try: myMQTTClient.publish( topic="home/helloworld", QoS=1, payload=publish_payload ) # 最多等待200ms确认发布成功 wait_start = time.time() while not publish_success and time.time() - wait_start < 0.2: myMQTTClient.loop() # 发布失败则存入兜底队列 if not publish_success: pending_send.extend(payload_list) except Exception as e: # 发布异常也存入兜底队列 pending_send.extend(payload_list) # 每10秒把兜底队列写入本地文件,避免程序崩溃丢失数据 if int(time.time()) % 10 == 0 and pending_send: try: with open(BACKUP_FILE, "w", encoding="utf-8") as f: json.dump(list(pending_send), f) except: pass
补充优化说明
- 上报数据增加窗口时间戳,下游AWS Lambda可以通过时间戳判断是否有重复上报或者漏上报的窗口
- 如果单窗口数据量过大,可以开启MQTT payload压缩,比如用gzip压缩后再发布,减少传输带宽占用
- 如果运行环境为Windows,
select对普通文件流支持有限,可以改用子线程单独读取遥测数据存入线程安全的queue.Queue,主循环每1秒从队列拉取所有累计数据批量上报即可
内容的提问来源于stack exchange,提问作者mehmetkocer
相关产品推荐
相关产品推荐

