Python多进程Manager:主进程变量更新、同步及客户端回调实现
技术需求
- 在使用
multiprocessing Manager的后台进程场景中,如何更新主进程中的变量? - 如何实现客户端与主进程中变量的同步?
- 当共享的
local_variable发生变化时,如何触发客户端侧的函数回调?
代码示例
共享Python类
import datetime class Sharing(): def __init__(self, value): self.c = value def getLocalTime(self): return datetime.datetime.now() def getC(self): return self.c def setC(self, value): self.c = value print(self.c)
带后台进程的服务端
from multiprocessing import Process from multiprocessing.managers import SyncManager import shared as shared import time class MyManager(SyncManager): pass # 后台进程共享服务 def start_bg_manager_server(local_variable): MyManager.register('share', local_variable.getC) manager = MyManager(address=('', 50000), authkey=b'test') s = manager.get_server() s.serve_forever() if __name__ == '__main__': local_variable = shared.Sharing(20) p = Process(target=start_bg_manager_server, args=(local_variable,)) p.start() print("DONE") while True: local_variable.setC(local_variable.getC() + 10) time.sleep(2)
服务端运行输出
e c:/Users/dev/Documents/poc_python/server.py DONE 30 40 50 60 70 80 90
简单客户端消费者
from multiprocessing.managers import SyncManager class MyManager(SyncManager): pass MyManager.register('share') m = MyManager(address=('127.0.0.1', 50000), authkey=b'test') m.connect() # print(m.share().setC(20)) print((m.share()))
客户端运行输出
PS C:\Users\dev\Documents\poc_python> python3 .\client.py 20
解决方案
1. 后台进程中更新主进程变量
当前代码仅注册了getC方法,且后台进程拿到的是主进程变量的副本,无法直接修改主进程数据。要实现跨进程更新,需将整个共享实例注册为远程可访问对象:
# 修改服务端的start_bg_manager_server函数 def start_bg_manager_server(local_variable): # 注册共享实例,通过lambda返回主进程的local_variable引用 MyManager.register('get_share', callable=lambda: local_variable) manager = MyManager(address=('', 50000), authkey=b'test') s = manager.get_server() s.serve_forever()
客户端通过get_share()获取实例的远程代理后,调用setC会直接作用于主进程的local_variable。
2. 客户端与主进程变量同步
通过远程代理定期获取最新值即可实现同步:
# 修改客户端代码 from multiprocessing.managers import SyncManager import time class MyManager(SyncManager): pass MyManager.register('get_share') m = MyManager(address=('127.0.0.1', 50000), authkey=b'test') m.connect() share_obj = m.get_share() while True: print(f"当前同步值: {share_obj.getC()}") time.sleep(2)
3. 变量变化触发客户端回调
Python原生SyncManager不支持主动推送,可通过两种方式实现:
方案1:轮询对比触发
客户端记录上次值,每次轮询对比,变化时触发回调:
from multiprocessing.managers import SyncManager import time class MyManager(SyncManager): pass MyManager.register('get_share') m = MyManager(address=('127.0.0.1', 50000), authkey=b'test') m.connect() share_obj = m.get_share() last_value = share_obj.getC() def on_value_change(new_value): print(f"值更新通知: {new_value}") while True: current_value = share_obj.getC() if current_value != last_value: on_value_change(current_value) last_value = current_value time.sleep(1)
方案2:消息队列主动推送
借助Queue实现服务端主动推送,变量变化时发送消息:
- 服务端修改:
from multiprocessing import Process, Queue from multiprocessing.managers import SyncManager import shared as shared import time class MyManager(SyncManager): pass def start_bg_manager_server(local_variable, update_queue): MyManager.register('get_share', callable=lambda: local_variable) MyManager.register('get_update_queue', callable=lambda: update_queue) manager = MyManager(address=('', 50000), authkey=b'test') s = manager.get_server() s.serve_forever() if __name__ == '__main__': local_variable = shared.Sharing(20) update_queue = Queue() p = Process(target=start_bg_manager_server, args=(local_variable, update_queue)) p.start() print("DONE") while True: new_value = local_variable.getC() + 10 local_variable.setC(new_value) update_queue.put(new_value) # 推送更新消息 time.sleep(2)
- 客户端修改:
from multiprocessing.managers import SyncManager class MyManager(SyncManager): pass MyManager.register('get_update_queue') m = MyManager(address=('127.0.0.1', 50000), authkey=b'test') m.connect() update_queue = m.get_update_queue() def on_value_change(new_value): print(f"收到更新通知: {new_value}") while True: new_value = update_queue.get() on_value_change(new_value)
内容的提问来源于stack exchange,提问作者yop038
相关产品推荐
相关产品推荐

