Multiprocessing Manager Dict跨进程无法更新问题求助
multiprocessing Manager字典跨进程共享数据失效问题解决
使用multiprocessing的Manager.dict()实现跨进程数据共享时,主进程始终无法获取子进程更新的数据,排查后发现问题出在子进程对共享字典的操作方式上。以下是可复现问题的最小示例代码:
from multiprocessing import Process, Manager import sys import time import random def run_pipeline(position, frame): while True: position = {1: random.randint(0, 100)} frame = {2: random.randint(0, 100)} print("pipeline running") time.sleep(5) def run_relay_controller(relay_status): while True: print("relay running") relay_status = {2: random.randint(0, 2)} time.sleep(5) def run_webserver(frame, relay_status): while True: print("webserver running") frame = {1: random.randint(0, 100)} relay_status = {1: random.randint(0, 2)} time.sleep(5) if __name__ == "__main__": manager = Manager() position = manager.dict() frame = manager.dict() relay_status = manager.dict() position_process = Process(target=run_pipeline, args=(position, frame)) relay_process = Process(target=run_relay_controller, args=(relay_status, )) server_process = Process(target=run_webserver, args=(frame, relay_status)) position_process.start() relay_process.start() server_process.start() while True: try: print(frame) print(relay_status) print(position) except KeyboardInterrupt: position_process.terminate() relay_process.terminate() server_process.terminate() break position_process.join() relay_process.join() server_process.join() sys.exit()
问题原因
子进程中使用position = {1: ...}这类赋值语句时,是将传入的共享Manager字典对象直接替换成了普通的Python字典。后续的修改操作都只针对这个局部普通字典,完全和原共享字典无关,导致主进程无法接收到更新。
解决方法
要修改共享字典的内容,而不是重新赋值覆盖原对象。可以通过两种方式实现:
- 直接给字典的键赋值:
position[1] = random.randint(0, 100) - 使用
update()方法批量更新:position.update({1: random.randint(0, 100)})
修改后的代码示例
from multiprocessing import Process, Manager import sys import time import random def run_pipeline(position, frame): while True: # 直接修改共享字典的键值 position[1] = random.randint(0, 100) frame[2] = random.randint(0, 100) print("pipeline running") time.sleep(5) def run_relay_controller(relay_status): while True: print("relay running") relay_status[2] = random.randint(0, 2) time.sleep(5) def run_webserver(frame, relay_status): while True: print("webserver running") frame[1] = random.randint(0, 100) relay_status[1] = random.randint(0, 2) time.sleep(5) if __name__ == "__main__": manager = Manager() position = manager.dict() frame = manager.dict() relay_status = manager.dict() position_process = Process(target=run_pipeline, args=(position, frame)) relay_process = Process(target=run_relay_controller, args=(relay_status, )) server_process = Process(target=run_webserver, args=(frame, relay_status)) position_process.start() relay_process.start() server_process.start() while True: try: print(frame) print(relay_status) print(position) except KeyboardInterrupt: position_process.terminate() relay_process.terminate() server_process.terminate() break position_process.join() relay_process.join() server_process.join() sys.exit()
运行修改后的代码,主进程就能正常获取到各子进程更新的共享字典数据了。
内容的提问来源于stack exchange,提问作者Moritz Pfennig
相关产品推荐
相关产品推荐

