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

在While True循环中用Multiprocessing实现PyMongo无阻塞上传的问题

解决串口数据多进程上传MongoDB的空文档与重复执行问题

问题根源

  1. 多进程内存隔离:multiprocessing的子进程与主进程内存空间完全独立,如果你在创建进程时RB_DataDict还未填充数据,子进程拿到的就是空字典的副本;即便后续主进程填充了数据,子进程也无法感知到。
  2. 进程仅初始化一次:如果你的代码只创建了一次Process对象并启动,自然只会执行一次上传操作,不会跟随循环重复运行。

修正后的代码示例

from multiprocessing import Process
import pymongo
import serial
import datetime

def upload_to_database(data):
    # 子进程内独立创建MongoDB连接(连接对象不能跨进程共享)
    client = pymongo.MongoClient("mongodb://localhost:27017/")
    db = client["sensor_db"]
    collection = db["serial_data"]
    collection.insert_one(data)
    client.close()

def main():
    # 初始化串口(根据实际端口、波特率修改)
    ser = serial.Serial('COM3', 9600, timeout=1)
    while True:
        RB_DataDict = {}
        # 读取串口数据并填充字典
        raw_data = ser.readline().decode('utf-8').strip()
        if raw_data:
            # 示例解析逻辑:假设数据格式为 "sensor1_val,sensor2_val"
            try:
                s1_val, s2_val = map(float, raw_data.split(','))
                RB_DataDict["sensor1"] = s1_val
                RB_DataDict["sensor2"] = s2_val
                RB_DataDict["upload_time"] = datetime.datetime.now()
            except ValueError:
                print("数据格式错误,跳过本次上传")
                continue
        
        # 数据非空时启动新进程上传
        if RB_DataDict:
            # 传递字典副本,确保子进程拿到的是当前完整数据
            upload_process = Process(target=upload_to_database, args=(RB_DataDict.copy(),))
            upload_process.start()
            # 不要调用join(),否则会阻塞主进程的串口读取循环

if __name__ == "__main__":
    main()

关键修改说明

  • 显式传递数据副本:用RB_DataDict.copy()作为参数传给子进程,确保子进程获取的是当前循环填充完成的完整数据,避免主进程后续修改(或下一次循环覆盖)影响子进程的上传内容。
  • 每次循环创建新进程:在每次数据填充完成后,重新实例化Process并启动,保证循环每执行一次就触发一次上传。
  • 避免全局变量依赖:让upload_to_database通过参数接收数据,而非依赖全局变量——多进程的全局变量是各自独立的副本,子进程无法共享主进程更新后的全局数据。
  • 独立创建MongoDB连接:MongoDB的客户端连接对象不能跨进程共享,必须在子进程内部单独创建连接。

优化建议(高频率数据场景)

如果串口数据频率很高,频繁创建销毁进程会有性能开销,建议改用进程池复用进程:

from multiprocessing import Pool

def main():
    ser = serial.Serial('COM3', 9600, timeout=1)
    # 初始化固定大小的进程池(根据CPU核心数或业务需求调整)
    with Pool(processes=3) as pool:
        while True:
            RB_DataDict = {}
            # 读取并填充数据逻辑...
            if RB_DataDict:
                # 异步提交任务到进程池,不阻塞主进程
                pool.apply_async(upload_to_database, args=(RB_DataDict.copy(),))

内容的提问来源于stack exchange,提问作者czar1249

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 06:35:24