如何解决Flask-SocketIO+Redis Pub/Sub多进程Nginx部署的重复消息问题
问题分析与解决方案
重复消息原因
你遇到的重复消息问题根源在于:
- 启动的两个Flask-SocketIO进程(5000、5001)各自启动了一个
listen子进程,两个子进程同时订阅Redis的my-channel频道。 - 当
publish.py向Redis频道发送一条消息时,两个listen子进程都会收到消息,各自调用socketio.emit('update', data)。 - 由于SocketIO配置了Redis消息队列,每次
emit都会将广播指令发送到Redis队列,两个主进程都会消费该指令并向各自的客户端广播。 - 最终连接到任意一个进程的客户端都会收到两次相同消息(一次来自自己连接的主进程处理第一个
listen进程的emit,一次处理第二个listen进程的emit)。
另外你的publish.py存在两处代码缺失:缺少json模块导入和Redis客户端初始化,需要补上才能正常运行。
解决方案
方案1:单全局监听进程(推荐)
不需要每个主进程都启动监听子进程,单独启动一个全局监听进程处理Redis消息,统一发送广播指令到SocketIO消息队列,避免重复触发广播。
步骤1:分离监听脚本
创建独立的redis_listener.py脚本:
import json import redis from flask_socketio import SocketIO REDIS_HOST = 'localhost' REDIS_PORT = 6379 REDIS_CHANNEL = 'my-channel' # 仅初始化SocketIO消息队列,无需Flask应用 socketio = SocketIO(message_queue=f'redis://{REDIS_HOST}:{REDIS_PORT}/') def listen(): redis_client = redis.Redis(host=REDIS_HOST, port=REDIS_PORT) pubsub = redis_client.pubsub() pubsub.subscribe(REDIS_CHANNEL) try: while True: message = pubsub.get_message() if message is not None and message['type'] == 'message': data = message['data'].decode('utf-8') data = json.loads(data) print(data) # 一次emit指令,所有SocketIO进程都会接收并向各自客户端广播 socketio.emit('update', data) finally: pubsub.close() if __name__ == '__main__': listen()
步骤2:修改app.py
移除启动listen子进程的代码:
from gevent import monkey monkey.patch_all() import sys from flask import Flask, render_template from flask_socketio import SocketIO, emit app = Flask(__name__) app.config['SECRET_KEY'] = 'secret!' REDIS_HOST = 'localhost' REDIS_PORT = 6379 socketio = SocketIO(app, message_queue=f'redis://{REDIS_HOST}:{REDIS_PORT}/') @app.route('/') def index(): return render_template('index.html') @socketio.on('connect') def handle_connect(): print('Client connected') emit('connected', {'data': 'Connected successfully!'}) @socketio.on('disconnect') def handle_disconnect(): print('Client disconnected') if __name__ == '__main__': port = int(sys.argv[1]) if len(sys.argv) > 1 else 5000 socketio.run(app, host='127.0.0.1', port=port)
步骤3:修正publish.py
补上缺失的代码:
import time import json import redis REDIS_HOST = 'localhost' REDIS_PORT = 6379 REDIS_CHANNEL = 'my-channel' redis_client = redis.Redis(host=REDIS_HOST, port=REDIS_PORT) while True: redis_client.publish( REDIS_CHANNEL, json.dumps(dict( msg='hello' )) ) time.sleep(3)
启动流程
- 启动全局监听进程:
python redis_listener.py - 启动两个SocketIO主进程:
python app.py 5000和python app.py 5001 - 启动消息发布脚本:
python publish.py
方案2:进程专属房间广播
如果必须保留每个主进程的监听逻辑,可以通过SocketIO房间特性,让每个进程只向自己连接的客户端广播。
修改app.py:
from gevent import monkey monkey.patch_all() import sys import json from flask import Flask, render_template from flask_socketio import SocketIO, emit, join_room import redis app = Flask(__name__) app.config['SECRET_KEY'] = 'secret!' REDIS_HOST = 'localhost' REDIS_PORT = 6379 REDIS_CHANNEL = 'my-channel' # 生成当前进程的专属房间标识 PROCESS_ROOM = f'process_{sys.argv[1]}' if len(sys.argv) > 1 else 'process_5000' socketio = SocketIO(app, message_queue=f'redis://{REDIS_HOST}:{REDIS_PORT}/') def listen(): redis_client = redis.Redis(host=REDIS_HOST, port=REDIS_PORT) pubsub = redis_client.pubsub() pubsub.subscribe(REDIS_CHANNEL) try: while True: message = pubsub.get_message() if message is not None and message['type'] == 'message': data = message['data'].decode('utf-8') data = json.loads(data) print(data) # 仅向当前进程的专属房间广播 socketio.emit('update', data, room=PROCESS_ROOM) finally: pubsub.close() @app.route('/') def index(): return render_template('index.html') @socketio.on('connect') def handle_connect(): print('Client connected') # 将客户端加入当前进程的专属房间 join_room(PROCESS_ROOM) emit('connected', {'data': 'Connected successfully!'}) @socketio.on('disconnect') def handle_disconnect(): print('Client disconnected') if __name__ == '__main__': from multiprocessing import Process process = Process(target=listen) process.start() port = int(sys.argv[1]) if len(sys.argv) > 1 else 5000 socketio.run(app, host='127.0.0.1', port=port)
此方案中,客户端连接时会被加入对应进程的专属房间,每个进程的监听逻辑只会向自己的房间广播消息,避免重复推送。
内容的提问来源于stack exchange,提问作者user21462504
相关产品推荐
相关产品推荐

