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

Python线程间数据共享问题:UDP数据分发给多线程实时处理

解决方案

针对你的需求,推荐两种高效的线程数据分发方案,解决全局变量冻结和队列数据丢失的问题:

方案1:多队列分发(适用于需要处理所有历史数据的场景)

为每个处理线程创建专属队列,消费者线程根据数据块ID将数据分发到对应队列,每个线程仅从自己的队列取数据,避免数据争抢。

实现代码:

import queue
import threading

# 为每个处理线程创建专属队列
queue_dict = {
    "id_1": queue.Queue(),
    "id_2": queue.Queue(),
    "id_3": queue.Queue()
}

def mlFunc(target_id):
    while True:
        # 从对应队列获取数据,阻塞等待新数据
        temp = queue_dict[target_id].get()
        canDo.publish(temp)
        # 标记任务完成(可选,用于队列跟踪)
        queue_dict[target_id].task_done()

def consumer(in_q):
    while True:
        canAr = canParser.b2can(in_q.get())
        # 假设canAr包含id字段和data字段
        data_id = canAr["id"]
        # 根据ID分发到对应队列
        if data_id in queue_dict:
            queue_dict[data_id].put(canAr["data"])

# 启动三个处理线程,分别绑定不同ID
threading.Thread(target=mlFunc, args=("id_1",), daemon=True).start()
threading.Thread(target=mlFunc, args=("id_2",), daemon=True).start()
threading.Thread(target=mlFunc, args=("id_3",), daemon=True).start()

方案2:线程安全字典+事件通知(适用于仅需处理最新数据的场景)

用threading.Lock保证字典线程安全,配合threading.Event通知对应线程有新数据。这种方式避免队列堆积,适合实时性要求高、不需要保留历史数据的场景。

实现代码:

import threading

# 存储最新数据的线程安全字典
latest_data = {}
# 每个ID对应一个事件,用于通知线程更新
event_dict = {
    "id_1": threading.Event(),
    "id_2": threading.Event(),
    "id_3": threading.Event()
}
data_lock = threading.Lock()

def mlFunc(target_id):
    while True:
        # 等待事件触发
        event_dict[target_id].wait()
        # 获取最新数据
        with data_lock:
            temp = latest_data.get(target_id)
            # 重置事件,等待下一次更新
            event_dict[target_id].clear()
        if temp:
            canDo.publish(temp)

def consumer(in_q):
    while True:
        canAr = canParser.b2can(in_q.get())
        data_id = canAr["id"]
        data = canAr["data"]
        if data_id in event_dict:
            with data_lock:
                latest_data[data_id] = data
            # 触发对应线程的事件
            event_dict[data_id].set()

# 启动线程
threading.Thread(target=mlFunc, args=("id_1",), daemon=True).start()
threading.Thread(target=mlFunc, args=("id_2",), daemon=True).start()
threading.Thread(target=mlFunc, args=("id_3",), daemon=True).start()

注意事项

  • 使用daemon=True标记线程为守护线程,避免程序退出时线程残留。
  • 方案1中如果数据量过大,需考虑队列的最大长度,避免内存溢出(可通过queue.Queue(maxsize=N)设置)。
  • 方案2中如果数据更新频率远高于线程处理速度,会导致旧数据被覆盖,需根据业务场景选择。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 02:12:22