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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 14:24:05