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

