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

Django+RabbitMQ多标签页并发连接数据库问题求助

多标签页独立管理RabbitMQ消费者的解决方案

问题背景

我在项目中用pika/RabbitMQ从网页向远程数据库发送SQL脚本,数据库处理后返回结果。需求是支持用户通过多标签页同时连接不同数据库,加速远程维护。但遇到两个核心问题:

  • Django的Session Storage是标签页共享的,修改Session会影响所有标签页,没法单独控制每个标签页的状态
  • 目前用UUID作为消费者ID来控制启停,但Local Storage会被新标签页覆盖,导致无法准确对应到目标消费者

现有消费者代码:

import pika
from comandos.connectionAMQP.connectionSettings import getConnectionRabbitMQ

class MessageReturn:
    def __init__(self, consumer_key) -> None:
        self.consumer_key = consumer_key
        self.message_body = None
        self.is_consuming = True

    def set_message_body(self, body):
        self.message_body = body

    def get_message_body(self):
        return self.message_body

    def stop_consuming(self):
        self.is_consuming = False

# 存储MessageReturn实例的列表
consumers = []

def create_consumer(consumer_key):
    global consumers
    consumer = MessageReturn(consumer_key)
    consumers.append(consumer)
    return consumer

def stop_consume(consumer_key):
    global consumers
    for consumer in consumers:
        if consumer.consumer_key == consumer_key and consumer.is_consuming:
            consumer.stop_consuming()
            print(f"Consumption stopped for consumer {consumer_key}.")
            break
    else:
        print(f"No active consumption to stop for consumer {consumer_key}.")

def consume(consumer_key):
    connection = pika.BlockingConnection(getConnectionRabbitMQ())
    channel = connection.channel()

    def callback(ch, method, properties, body):
        for consumer in consumers:
            if consumer.consumer_key == consumer_key:
                consumer.set_message_body(body.decode('utf-8'))
                consumer.stop_consuming()

    global consumers
    state = create_consumer(consumer_key)

    channel.basic_consume(queue='Response_BKP_00788556-af71-4c99-a8c8-1fde8841e1ae', on_message_callback=callback, auto_ack=True)

    print(f'[*] Waiting for messages for consumer {consumer_key}. To exit press CTRL+C')

    while state.is_consuming:
        connection.process_data_events()

    connection.close()
    return state.get_message_body()

可行解决方案

1. 客户端生成标签页独立ID(快速落地方案)

  • 每个标签页加载时,用浏览器的crypto.randomUUID()生成唯一ID,存在sessionStorage里(sessionStorage是标签页隔离的,关闭标签自动销毁,不会被其他标签覆盖)
  • 每次和后端交互(启动/停止消费者)时,把这个ID作为参数传给后端
  • 后端用字典替代全局列表存储消费者,以ID为key,直接通过ID定位目标消费者,避免遍历查找的低效

2. WebSocket双向通信(优雅长连接方案)

  • 每个标签页打开时,和后端建立独立WebSocket连接,每个连接对应一个RabbitMQ消费者
  • 数据库返回结果后,后端直接通过对应WebSocket推给标签页,无需轮询或依赖Session
  • 标签页关闭时,WebSocket连接断开,后端可直接销毁对应消费者,避免资源浪费
  • Django中可以用channels库实现WebSocket,配合ASGI服务器(如Daphne)运行

3. 优化后端消费者管理逻辑

  • 把全局consumers列表改成字典,key为消费者ID,value存储MessageReturn实例+RabbitMQ连接/通道,方便快速定位和资源清理
  • 增加超时清理机制,定期检查长时间无响应的消费者,自动销毁释放资源
  • 修改后的核心代码示例:
# 替换全局列表为字典,存储消费者实例及连接信息
consumers = {}

def create_consumer(consumer_key):
    consumer = MessageReturn(consumer_key)
    consumers[consumer_key] = {
        "instance": consumer,
        "connection": None,
        "channel": None
    }
    return consumer

def stop_consume(consumer_key):
    consumer_info = consumers.get(consumer_key)
    if not consumer_info:
        print(f"No active consumption to stop for consumer {consumer_key}.")
        return
    consumer = consumer_info["instance"]
    if consumer.is_consuming:
        consumer.stop_consuming()
        # 主动关闭连接释放资源
        if consumer_info["connection"]:
            consumer_info["connection"].close()
        del consumers[consumer_key]
        print(f"Consumption stopped for consumer {consumer_key}.")

def consume(consumer_key):
    connection = pika.BlockingConnection(getConnectionRabbitMQ())
    channel = connection.channel()
    # 保存连接和通道到字典
    consumers[consumer_key]["connection"] = connection
    consumers[consumer_key]["channel"] = channel

    def callback(ch, method, properties, body):
        consumer_info = consumers.get(consumer_key)
        if not consumer_info:
            return
        consumer = consumer_info["instance"]
        if consumer.is_consuming:
            consumer.set_message_body(body.decode('utf-8'))
            consumer.stop_consuming()

    state = create_consumer(consumer_key)

    channel.basic_consume(queue='Response_BKP_00788556-af71-4c99-a8c8-1fde8841e1ae', on_message_callback=callback, auto_ack=True)

    print(f'[*] Waiting for messages for consumer {consumer_key}. To exit press CTRL+C')

    while state.is_consuming:
        connection.process_data_events()

    # 清理资源
    connection.close()
    if consumer_key in consumers:
        del consumers[consumer_key]
    return state.get_message_body()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 09:34:55