Python3中如何在独立启动的进程之间共享Queue队列?
Python3 独立进程间共享队列解决方案(基于BaseManager)
整体逻辑
BaseManager通过本地套接字/网络通信实现跨独立进程的对象共享,不需要进程间有父子关系,完全符合按需独立启动进程的需求。队列实例统一维护在channel端,station端通过代理对象操作队列,默认基于pickle序列化,直接支持BitArray对象传输。
1. 双端共通定义
channel端和station端都需要保证以下配置完全一致:
from multiprocessing.managers import BaseManager from queue import Queue from bitarray import BitArray # 服务配置,双端必须完全相同 MANAGER_ADDR = ('127.0.0.1', 50000) # 跨机器部署改为对应IP AUTH_KEY = b'your_custom_secret_key' # 自定义密钥,bytes类型
2. Channel端(服务端,先启动)
负责维护所有队列实例,启动manager服务供station连接:
# 初始化全局队列 channel_queue = Queue() # station写,channel读 station_queues = dict() # 存储各station专属队列,channel写,对应station读 # 定义暴露给station的方法 def get_channel_queue(): return channel_queue def register_station_queue(station_id: str): if station_id not in station_queues: station_queues[station_id] = Queue() def get_station_queue(station_id: str): return station_queues.get(station_id) # 注册方法到manager BaseManager.register('get_channel_queue', callable=get_channel_queue) BaseManager.register('register_station_queue', callable=register_station_queue) BaseManager.register('get_station_queue', callable=get_station_queue) # 启动服务 if __name__ == '__main__': manager = BaseManager(address=MANAGER_ADDR, authkey=AUTH_KEY) server = manager.get_server() print(f"Channel服务启动,监听地址{MANAGER_ADDR}") server.serve_forever()
Channel端读写逻辑
- 读取station上报数据:直接调用
channel_queue.get()即可 - 向指定station下发
BitArray数据:station_queues[目标station_id].put(bit_array_obj)
3. Station端(客户端,后启动)
不需要维护队列实例,直接连接channel服务获取队列代理对象操作:
# 注册方法(不需要传callable,客户端使用服务端的实现) BaseManager.register('get_channel_queue') BaseManager.register('register_station_queue') BaseManager.register('get_station_queue') if __name__ == '__main__': # 连接channel服务 manager = BaseManager(address=MANAGER_ADDR, authkey=AUTH_KEY) manager.connect() # 替换为当前station的唯一ID,不可重复 STATION_ID = "station_001" # 注册当前station的专属队列 manager.register_station_queue(STATION_ID) # 获取队列代理对象 channel_queue = manager.get_channel_queue() my_queue = manager.get_station_queue(STATION_ID) # 读写示例 # 向channel上报数据 channel_queue.put(f"来自{STATION_ID}的消息") # 读取channel下发的BitArray数据 recv_bit_array = my_queue.get() print(f"收到数据:{recv_bit_array}")
注意事项
- 必须先启动channel服务,再启动station进程
- station_id需要全局唯一,避免队列覆盖
- 自定义对象如果默认序列化失败,自行实现
__reduce__方法适配pickle即可 - 队列代理对象的
put/get/put_nowait/get_nowait等方法和原生Queue用法完全一致
内容的提问来源于stack exchange,提问作者Mahek Shamsukha
相关产品推荐
相关产品推荐

