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

如何解决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)

启动流程

  1. 启动全局监听进程:python redis_listener.py
  2. 启动两个SocketIO主进程:python app.py 5000 和 python app.py 5001
  3. 启动消息发布脚本: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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 09:37:04