如何从Python守护线程向多线程定期分发dict()?
多线程/进程共享实时更新字典的解决方案
嘿,你的场景其实挺典型的——一个生产者(守护线程更新数据)对应多个消费者(1-15个线程/进程定期获取数据),原代码的问题在于用了普通的Lock,它一次只能唤醒一个等待的线程,所以没法实现一对多的分发。下面给你几个简便且可靠的实现方式,全都是基于Python标准库的:
一、线程场景:用threading.Condition实现一对多通知
Condition是普通Lock的增强版,它支持通知所有等待的线程,完美适配你的一个更新线程对应多个消费线程的需求。核心逻辑是:
- 更新线程每次完成字典更新后,主动通知所有等待的消费线程
- 消费线程可以选择等待更新通知,或者定期主动拉取(兼顾实时性和定期获取的需求)
示例代码:
import time import threading as t from datetime import datetime def updater(mydict, period, cond): while True: # 模拟WebSocket获取实时数据 current_sec = int(str(datetime.now())[17:19]) # 按周期更新字典 if current_sec % period == 0: with cond: mydict['data'] = current_sec cond.notify_all() # 通知所有等待的消费线程 time.sleep(1) def calculator(mydict, period, cond, thread_id): while True: with cond: # 等待更新通知,最多等待period秒(超时自动拉取最新数据) cond.wait(timeout=period) # 复制字典避免后续更新影响当前计算 data = mydict.copy() # 模拟业务计算逻辑 print(f"线程{thread_id} - 当前时间: {int(str(datetime.now())[17:19])}, 获取数据: {data}") if __name__ == '__main__': shared_dict = {} period = 5 cond = t.Condition() # 启动守护更新线程(主进程退出时自动终止) updater_thread = t.Thread(target=updater, args=(shared_dict, period, cond), daemon=True) updater_thread.start() # 启动多个消费线程(示例为3个) for i in range(3): calc_thread = t.Thread(target=calculator, args=(shared_dict, period, cond, i+1)) calc_thread.start() # 保持主进程运行 updater_thread.join()
这个实现的优势:
- 所有消费线程都能及时收到数据更新通知
- 用
with cond自动管理锁的获取和释放,避免手动操作锁的遗漏 - 通过
wait(timeout)同时支持“实时响应更新”和“定期拉取”两种模式
二、进程场景:用multiprocessing共享字典+Condition
如果你的消费者是进程(而非线程),因为进程间内存不共享,普通字典无法直接复用,需要用multiprocessing.Manager创建跨进程共享的字典,搭配multiprocessing.Condition实现通知:
示例代码:
import time import multiprocessing as mp from datetime import datetime def updater(mydict, period, cond): while True: current_sec = int(str(datetime.now())[17:19]) if current_sec % period == 0: with cond: mydict['data'] = current_sec cond.notify_all() time.sleep(1) def calculator(mydict, period, cond, proc_id): while True: with cond: cond.wait(timeout=period) data = mydict.copy() print(f"进程{proc_id} - 当前时间: {int(str(datetime.now())[17:19])}, 获取数据: {data}") if __name__ == '__main__': with mp.Manager() as manager: shared_dict = manager.dict() period = 5 cond = manager.Condition() # 启动守护更新进程 updater_proc = mp.Process(target=updater, args=(shared_dict, period, cond), daemon=True) updater_proc.start() # 启动多个消费进程 procs = [] for i in range(3): proc = mp.Process(target=calculator, args=(shared_dict, period, cond, i+1)) procs.append(proc) proc.start() for proc in procs: proc.join()
三、更简便的“定期快照”方案(无需通知)
如果你的消费者不需要实时响应更新,只是定期获取最新数据,那可以不用Condition,直接给共享字典加个普通锁,消费者定期加锁读取即可:
线程场景示例:
import time import threading as t from datetime import datetime def updater(mydict, period, lock): while True: current_sec = int(str(datetime.now())[17:19]) if current_sec % period == 0: with lock: mydict['data'] = current_sec time.sleep(1) def calculator(mydict, period, lock, thread_id): while True: time.sleep(period) with lock: data = mydict.copy() print(f"线程{thread_id} - 当前时间: {int(str(datetime.now())[17:19])}, 获取数据: {data}") if __name__ == '__main__': shared_dict = {} period = 5 lock = t.Lock() updater_thread = t.Thread(target=updater, args=(shared_dict, period, lock), daemon=True) updater_thread.start() for i in range(3): calc_thread = t.Thread(target=calculator, args=(shared_dict, period, lock, i+1)) calc_thread.start() updater_thread.join()
这个方案逻辑最简单,适合对实时性要求不高的场景,缺点是消费者可能会错过两次间隔之间的临时更新,但刚好匹配你“定期获取”的需求。
总结选择建议:
- 如果是线程场景且需要实时通知:优先用
threading.Condition - 如果是进程场景且需要实时通知:用
multiprocessing.Manager的共享字典+Condition - 如果只需要定期获取最新数据:不管线程还是进程,用锁保护共享字典即可
内容的提问来源于stack exchange,提问作者back_pockets
相关产品推荐
相关产品推荐

