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
相关产品推荐
相关产品推荐

