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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 11:45:03