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

能否编写Python脚本为多客户批量创建预配置的RabbitMQ队列?

RabbitMQ批量创建客户队列Python实现方案

前置依赖

  • 安装Python RabbitMQ客户端库:pip install pika
  • 提前梳理预定义的公共队列配置、交换机信息,以及各客户的差异化参数(队列名称、routing key等)

核心实现代码

import pika
from pika.exceptions import AMQPConnectionError

# -------------------------- 公共统一配置 --------------------------
RABBITMQ_HOST = "你的RabbitMQ服务地址"
RABBITMQ_PORT = 5672
RABBITMQ_USER = "服务用户名"
RABBITMQ_PWD = "服务密码"
VHOST = "业务使用的vhost,默认值为/"
# 队列公共基础配置,所有客户默认复用该配置
COMMON_QUEUE_ARGS = {
    "durable": True,  # 队列持久化
    "exclusive": False, # 非排他队列,允许多个客户端连接
    "auto_delete": False, # 消费完不自动删除队列
    "arguments": {"x-message-ttl": 86400000} # 公共额外参数,示例为消息1天过期,可按需调整
}
# 提前创建的业务统一预定义交换机
PREDEFINED_EXCHANGE = "业务公用交换机名称"
EXCHANGE_TYPE = "direct" # 按需替换为topic/headers等实际交换机类型

# -------------------------- 客户差异化配置 --------------------------
# 仅需填写每个客户的专属参数,支持单独覆盖公共配置
CUSTOMER_CONFIGS = [
    {
        "customer_id": "customer_001",
        "queue_name": "cus_001_biz_queue",
        "routing_key": "biz.cus001"
    },
    {
        "customer_id": "customer_002",
        "queue_name": "cus_002_biz_queue",
        "routing_key": "biz.cus002",
        # 个别客户需要特殊配置时新增该字段,优先级高于公共配置
        "override_args": {
            "arguments": {"x-message-ttl": 3600000}
        }
    }
]

def batch_create_customer_queues():
    # 建立RabbitMQ连接
    credentials = pika.PlainCredentials(RABBITMQ_USER, RABBITMQ_PWD)
    try:
        connection = pika.BlockingConnection(pika.ConnectionParameters(
            host=RABBITMQ_HOST,
            port=RABBITMQ_PORT,
            virtual_host=VHOST,
            credentials=credentials
        ))
    except AMQPConnectionError as e:
        print(f"RabbitMQ连接失败:{str(e)}")
        return
    channel = connection.channel()

    # 确保预定义交换机存在
    channel.exchange_declare(exchange=PREDEFINED_EXCHANGE, exchange_type=EXCHANGE_TYPE, durable=True)

    # 遍历客户配置批量创建队列、绑定routing key
    for config in CUSTOMER_CONFIGS:
        # 合并配置:客户自定义配置优先于公共配置
        queue_config = COMMON_QUEUE_ARGS.copy()
        if "override_args" in config:
            queue_config.update(config["override_args"])
        
        # 声明队列
        channel.queue_declare(
            queue=config["queue_name"],
            durable=queue_config["durable"],
            exclusive=queue_config["exclusive"],
            auto_delete=queue_config["auto_delete"],
            arguments=queue_config["arguments"]
        )
        # 绑定队列到指定交换机
        channel.queue_bind(
            queue=config["queue_name"],
            exchange=PREDEFINED_EXCHANGE,
            routing_key=config["routing_key"]
        )
        print(f"客户{config['customer_id']}配置完成:队列名{config['queue_name']},routing key{config['routing_key']}")

    connection.close()
    print("全部客户队列批量创建、绑定完成")

if __name__ == "__main__":
    batch_create_customer_queues()

使用说明

  • 先修改脚本顶部的RabbitMQ连接参数、公共配置、交换机信息为实际业务值
  • 客户配置列表CUSTOMER_CONFIGS可以对接内部客户管理系统,直接从数据库导出参数赋值即可,无需手动录入
  • 如需支持批量修改、删除队列,只需修改遍历逻辑中的对应pika方法即可

注意事项

  • 建议在RabbitMQ同内网机器运行脚本,避免连接超时问题
  • 首次运行可先注释批量逻辑,仅测试单个客户的配置是否符合预期,验证通过后再全量执行
  • 脚本基于AMQP协议实现,兼容所有主流版本的RabbitMQ

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 06:57:00