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

Python消费RabbitMQ后SocketIO推送Vue2失败问题排查

问题原因

核心问题是Python标准线程与eventlet协程上下文不兼容:

  • SocketIO服务基于eventlet的协程模型运行,而你用threading.Thread启动的RabbitMQ消费线程是标准OS线程,在这个线程中调用SocketIO的emit方法时,无法进入eventlet的异步IO上下文,导致消息无法正确推送到前端。
  • 前端主动触发的random_number事件能正常响应,是因为该回调运行在eventlet的协程上下文内,和SocketIO服务的运行环境一致。

解决方案

1. 添加eventlet猴子补丁

在代码最开头添加补丁,确保所有IO操作(包括pika的网络连接)适配eventlet的协程模型:

import eventlet
eventlet.monkey_patch()  # 必须放在所有其他import之前

2. 用eventlet协程替代标准线程

将启动RabbitMQ消费的标准线程替换为eventlet协程,确保回调函数运行在eventlet上下文内:
修改main函数中的线程创建逻辑:

def main():
    global threads
    load_dotenv()
    RABBITMQ_USER = os.getenv('RABBITMQ_USER')
    RABBITMQ_PW = os.getenv('RABBITMQ_PW')
    RABBITMQ_IP = os.getenv('RABBITMQ_IP')

    socketio_handler = SocketIOHandler()
    # 替换标准线程为eventlet协程
    rabbitmq_greenlet = eventlet.spawn(start_rabbitmq_handler, socketio_handler, RABBITMQ_USER, RABBITMQ_PW, RABBITMQ_IP)
    threads.append(rabbitmq_greenlet)
    socketio_handler.run()

3. 优化RabbitMQ回调中的emit调用(可选)

如果仍有问题,用start_background_task确保emit操作在eventlet协程中执行:

def start_rabbitmq_handler(socketio_handler, RABBITMQ_USER, RABBITMQ_PW, RABBITMQ_IP):
    def callback(ch, method, properties, body):
        logging.info('rabbitmq handler')
        # 用start_background_task包装emit操作
        socketio_handler.sio.start_background_task(
            socketio_handler.emit, 
            'response_random_number', 
            { 'number': random.randint(0,10)}
        )

    with RabbitMQHandler(RABBITMQ_USER, RABBITMQ_PW, RABBITMQ_IP) as rabbitmq_handler:
        rabbitmq_handler.run(callback=callback)

修改后的完整代码

import eventlet
eventlet.monkey_patch()  # 必须放在最开头

import random
import socketio
import sys
import os
import uuid
import pika
from dotenv import load_dotenv
import logging

class RabbitMQHandler():
    def __init__(self, RABBITMQ_USER, RABBITMQ_PW, RABBITMQ_IP):
        self.queue_name = 'myqueue'
        self.exchange_name = 'myqueue'
        credentials = pika.PlainCredentials(RABBITMQ_USER, RABBITMQ_PW)
        self.connection = pika.BlockingConnection(pika.ConnectionParameters(RABBITMQ_IP, 5672, '/', credentials))
        self.channel = self.connection.channel()

        self.channel.queue_declare(queue=self.queue_name)
        self.channel.exchange_declare(exchange=self.exchange_name, exchange_type='fanout')
        self.channel.queue_bind(exchange=self.exchange_name, queue=self.queue_name)

    def __enter__(self):
        return self

    def __exit__(self, exc_type, exc_value, traceback):
        self.connection.close()

    def run(self, callback):
        logging.info('start consuming messages...')
        self.channel.basic_consume(queue=self.queue_name,auto_ack=True, on_message_callback=callback)
        self.channel.start_consuming()


class SocketIOHandler():
    def __init__(self):
        self.id = str(uuid.uuid4())
        # create a Socket.IO server
        self.sio = socketio.Server(async_mode='eventlet', cors_allowed_origins='*')
        # wrap with a WSGI application
        self.app = socketio.WSGIApp(self.sio)

        self.sio.on('connect_to_backend', self.handle_connect)
        self.sio.on('random_number', self.handle_random_number)

    def handle_connect(self, sid, msg):
        logging.info('new socket io connection from sid: {}'.format(sid))
        self.emit('connect_success', {
            'success': True,
        }, sid=sid)  # 指定sid只推送给当前连接的客户端

    def handle_random_number(self, sid, msg):
        logging.info('handle_random_number from sid: {}'.format(sid))
        self.emit('response_random_number', { 'number': random.randint(0,10)}, sid=sid)

    def emit(self, event, msg, sid=None):
        logging.info('socket server: {}'.format(self.id))
        logging.info('sending event: "{}"'.format(event))
        # 如果指定sid则推送给单个客户端,否则广播
        if sid:
            self.sio.emit(event, msg, room=sid)
        else:
            self.sio.emit(event, msg)
        logging.info('sent event: "{}"'.format(event))

    def run(self):
        logging.info('start web socket on port 8765...')
        eventlet.wsgi.server(eventlet.listen(('', 8765)), self.app)

def start_rabbitmq_handler(socketio_handler, RABBITMQ_USER, RABBITMQ_PW, RABBITMQ_IP):
    def callback(ch, method, properties, body):
        logging.info('received rabbitmq message')
        # 用start_background_task确保emit在eventlet协程中执行
        socketio_handler.sio.start_background_task(
            socketio_handler.emit, 
            'response_random_number', 
            { 'number': random.randint(0,10)}
        )

    with RabbitMQHandler(RABBITMQ_USER, RABBITMQ_PW, RABBITMQ_IP) as rabbitmq_handler:
        rabbitmq_handler.run(callback=callback)


threads = []

def main():
    global threads
    load_dotenv()
    RABBITMQ_USER = os.getenv('RABBITMQ_USER')
    RABBITMQ_PW = os.getenv('RABBITMQ_PW')
    RABBITMQ_IP = os.getenv('RABBITMQ_IP')

    socketio_handler = SocketIOHandler()
    # 使用eventlet协程启动RabbitMQ消费
    rabbitmq_greenlet = eventlet.spawn(start_rabbitmq_handler, socketio_handler, RABBITMQ_USER, RABBITMQ_PW, RABBITMQ_IP)
    threads.append(rabbitmq_greenlet)
    socketio_handler.run()

if __name__ == '__main__':
    try:
        logging.basicConfig(level=logging.INFO)
        logging.getLogger("pika").propagate = False
        main()
    except KeyboardInterrupt:
        try:
            for t in threads:
                t.kill()  # eventlet协程用kill终止
            sys.exit(0)
        except Exception:
            sys.exit(0)
    except SystemExit:
        for t in threads:
            t.kill()
        os._exit(0)

额外优化点

  • 在emit方法中添加sid参数,支持推送给单个客户端或广播,更灵活。
  • 终止eventlet协程要用kill()方法,而非标准线程的exit()。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 02:20:39