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

Python 3中实现扇出型进程间通信应选用什么模块?

嘿,这个需求我熟!刚好之前帮朋友做过类似的进程间消息推送场景,给你两个靠谱的实现方案,完全满足你的要求:进程A持续发消息,任意B进程随时启停,订阅后只收后续的消息,全程不用B给A发请求,自动推送~

方案一:基于Redis的Pub/Sub实现

Redis的发布/订阅(Pub/Sub)机制简直是为这种场景量身定做的——轻量、易集成,不管你的进程在同一台机器还是跨机器都能搞定,而且天生支持“订阅后仅接收新消息”的逻辑。

1. 准备依赖

先装Redis的Python客户端,同时确保你的机器上运行着Redis服务(本地或远程都可以):

pip install redis

2. 进程A(消息生产者)代码

这个进程会持续生成消息,发布到指定频道:

import redis
import time
import random

def producer():
    # 连接Redis服务
    r = redis.Redis(host='localhost', port=6379, db=0)
    # 指定消息发布的频道名称
    feed_channel = 'my_app_feed'
    
    print("进程A启动,开始发布消息...")
    msg_counter = 0
    while True:
        # 模拟生成业务消息内容
        message = f"Feed消息 #{msg_counter} - 随机内容: {random.randint(1, 1000)}"
        # 发布消息到频道
        r.publish(feed_channel, message)
        print(f"进程A已发布: {message}")
        msg_counter += 1
        time.sleep(1)  # 每秒发一条,可根据实际需求调整频率

if __name__ == "__main__":
    producer()

3. 进程B(消息消费者)代码

启动后自动订阅频道,只接收订阅之后发布的消息,关闭进程就停止订阅:

import redis
import sys

def consumer(consumer_name):
    # 连接Redis服务
    r = redis.Redis(host='localhost', port=6379, db=0)
    # 创建订阅对象
    pubsub = r.pubsub()
    # 订阅指定的消息频道
    pubsub.subscribe('my_app_feed')
    
    print(f"进程B[{consumer_name}]启动,开始订阅消息...")
    # 循环接收推送的消息
    for msg in pubsub.listen():
        # 过滤掉订阅确认的系统消息,只处理实际业务消息
        if msg['type'] == 'message':
            content = msg['data'].decode('utf-8')
            print(f"进程B[{consumer_name}]收到消息: {content}")

if __name__ == "__main__":
    # 给每个B进程指定唯一标识,方便区分不同消费者
    consumer_id = sys.argv[1] if len(sys.argv) > 1 else "默认消费者"
    consumer(consumer_id)

4. 运行方式

  • 先启动Redis服务(如果还没运行的话)
  • 启动进程A:python producer.py
  • 启动任意多个进程B:python consumer.py 消费者1、python consumer.py 消费者2,随时关闭再启动,新启动的B只会收到之后的消息,完全符合你的需求。

方案二:纯Python ZeroMQ实现(无外部依赖)

如果不想依赖Redis,也可以用ZeroMQ这个轻量消息库,它的Pub/Sub模式同样支持自动推送,而且是纯Python生态的方案。

1. 安装依赖

pip install pyzmq

2. 进程A(消息生产者)代码

import zmq
import time
import random

def producer():
    context = zmq.Context()
    socket = context.socket(zmq.PUB)
    # 绑定到本地5555端口,供B进程连接
    socket.bind("tcp://*:5555")
    
    print("进程A启动,开始发布消息...")
    msg_counter = 0
    while True:
        message = f"Feed消息 #{msg_counter} - 随机内容: {random.randint(1, 1000)}"
        socket.send_string(message)
        print(f"进程A已发布: {message}")
        msg_counter += 1
        time.sleep(1)

if __name__ == "__main__":
    producer()

3. 进程B(消息消费者)代码

import zmq
import sys

def consumer(consumer_name):
    context = zmq.Context()
    socket = context.socket(zmq.SUB)
    # 连接到进程A的发布端口
    socket.connect("tcp://localhost:5555")
    # 设置订阅规则,空字符串表示接收所有消息
    socket.setsockopt_string(zmq.SUBSCRIBE, "")
    
    print(f"进程B[{consumer_name}]启动,开始订阅消息...")
    while True:
        message = socket.recv_string()
        print(f"进程B[{consumer_name}]收到消息: {message}")

if __name__ == "__main__":
    consumer_id = sys.argv[1] if len(sys.argv) > 1 else "默认消费者"
    consumer(consumer_id)

关键注意事项
  • 不管用哪种方案,新启动的B进程只会收到订阅之后发布的消息,完全符合你的需求(如果需要回溯历史消息,Redis的Stream或者ZeroMQ的持久化功能可以实现,但你没提这个需求,就不多展开了)。
  • 多个B进程可以同时运行,消息会被推送给所有在线的订阅者。
  • 进程B随时启停都不会影响进程A的运行,A只管持续发布消息即可。

内容的提问来源于stack exchange,提问作者Roy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 08:58:31