欢迎光临
我们一直在努力

Ubuntu 系统下 AMQP 协议 RabbitMQ服务器部署

一、安装RabbitMQ及其配置

步骤1:安装RabbitMQ

# 更新系统
sudo apt update
sudo apt upgrade -y

# 安装Erlang和RabbitMQ
sudo apt install -y erlang
sudo apt install -y rabbitmq-server

# 启动并设置开机自启
sudo systemctl start rabbitmq-server
sudo systemctl enable rabbitmq-server

# 检查状态
sudo systemctl status rabbitmq-server

步骤2:安全配置(设置账户、密码 、可以正常接收局域网消息)

# 创建新用户(替换为你自己的用户名和密码)
sudo rabbitmqctl add_user mi_user YourStrongPassword123!

# 设置管理员权限(可选,根据需要)
sudo rabbitmqctl set_user_tags mi_user administrator

# 设置权限
sudo rabbitmqctl set_permissions -p / mi_user ".*" ".*" ".*"

# 删除默认的guest用户(必须!)
sudo rabbitmqctl delete_user guest

# 启用管理插件(可选,方便监控)
sudo rabbitmq-plugins enable rabbitmq_management

# 配置监听所有网络接口
echo 'listeners.tcp.default = 5672' | sudo tee -a /etc/rabbitmq/rabbitmq.conf
echo 'management.tcp.ip = 0.0.0.0' | sudo tee -a /etc/rabbitmq/rabbitmq.conf

# 重启服务使配置生效
sudo systemctl restart rabbitmq-server

# 开放防火墙
sudo ufw allow 5672/tcp
sudo ufw allow 15672/tcp

步骤3:验证安装

# 查看RabbitMQ状态
sudo rabbitmqctl status

# 查看监听端口
sudo netstat -tlnp | grep 5672
#查看所有用户
sudo rabbitmqctl list_users

步骤4:安装 pika 1.3.2 版本并使用阿里云镜像源的命令

pip install pika==1.3.2 -i https://mirrors.aliyun.com/pypi/simple/

二、Python程序(极简版)

  • set.py:运行一次,创建队列
  • producer.py:发送消息到队列
  • consumer.py:监听队列,接收消息
  • 1. set.py

    #!/usr/bin/env python3
    # -*- coding: utf-8 -*-
    """
    setup.py – RabbitMQ 交换机版基础设施配置脚本

    功能说明:
    1. 连接到 RabbitMQ Broker
    2. 创建交换机(持久化)
    3. 创建队列(持久化)
    4. 绑定队列到交换机(设置路由键)
    5. 提供运行提示信息

    作者:vb200811
    日期:2026-01-02
    """

    import pika

    # ======================== 配置参数 ========================
    HOST = 'localhost' # RabbitMQ 服务器的 IP 地址
    # 如果在同一台服务器,用 localhost;如果在其他设备,改成服务器 IP
    PORT = 5672 # RabbitMQ 默认 AMQP 端口(5672)
    USER = 'mi_user' # 用户名
    PASSWORD = 'YourPassword123' # 密码

    try:
    # 1. 创建连接(BlockingConnection 表示同步阻塞模式)
    conn = pika.BlockingConnection(
    pika.ConnectionParameters(
    host=HOST,
    port=PORT,
    credentials=pika.PlainCredentials(USER, PASSWORD) # 用户名密码认证
    )
    )

    # 2. 创建一个信道(Channel),所有操作都在信道中进行
    ch = conn.channel()
    print(f"✅ 连接成功 {HOST}:{PORT}")

    # 3. 创建交换机
    # exchange: 交换机名称
    # exchange_type: 交换机类型(direct/topic/fanout)
    # durable: 是否持久化(True 表示 RabbitMQ 重启后交换机仍然存在)
    ch.exchange_declare(
    exchange='my_exchange', # 交换机名称
    exchange_type='direct', # 类型:direct(直接匹配路由键)
    durable=True # 持久化
    )
    print("✅ 交换机 'my_exchange' 已创建")

    # 4. 创建队列
    # queue: 队列名称
    # durable: 是否持久化(True 表示 RabbitMQ 重启后队列仍然存在)
    ch.queue_declare(queue='high_queue', durable=True) # 高优先级队列
    ch.queue_declare(queue='low_queue', durable=True) # 低优先级队列
    print("✅ 队列已创建")

    # 5. 绑定队列到交换机
    # exchange: 交换机名称
    # queue: 队列名称
    # routing_key: 路由键(交换机根据路由键将消息路由到对应的队列)
    ch.queue_bind(
    exchange='my_exchange',
    queue='high_queue',
    routing_key='high' # 路由键 high -> high_queue
    )

    ch.queue_bind(
    exchange='my_exchange',
    queue='low_queue',
    routing_key='low' # 路由键 low -> low_queue
    )
    print("✅ 绑定完成")

    # 6. 关闭连接
    conn.close()
    print("\\n✅ 基础设施配置完成!")
    print("现在可以运行:")
    print(" python consumer.py # 启动消费者")
    print(" python producer.py # 发送消息")

    except Exception as e:
    # 捕获并处理异常
    print(f"❌ 配置失败: {e}")

    可运行下面命令查看配置:

    # 查看交换机
    sudo rabbitmqctl list_exchanges

    # 查看绑定
    sudo rabbitmqctl list_bindings

    2.producer.py

    #!/usr/bin/env python3
    # -*- coding: utf-8 -*-
    """
    producer.py – RabbitMQ 生产者示例(持久化 + 生产者确认模式)

    功能说明:
    1. 连接到 RabbitMQ Broker
    2. 开启生产者确认模式(Publisher Confirm)
    3. 发送持久化消息到交换机
    4. 处理消息路由失败等异常情况

    作者:vb200811
    日期:2026-01-02
    """

    import pika
    import json

    # ======================== 配置参数 ========================
    HOST = '192.168.1.100' # RabbitMQ 服务器的 IP 地址
    PORT = 5672 # RabbitMQ 默认 AMQP 端口(5672)
    USER = 'mi_user' # 用户名
    PASSWORD = 'YourPassword123' # 密码

    # ======================== 发送消息函数 ========================
    def send_message():
    """
    发送消息到 RabbitMQ 交换机,并使用生产者确认模式确保消息已被 Broker 接收并持久化。
    """
    # 1. 创建连接(BlockingConnection 表示同步阻塞模式)
    conn = pika.BlockingConnection(
    pika.ConnectionParameters(
    host=HOST,
    port=PORT,
    credentials=pika.PlainCredentials(USER, PASSWORD) # 用户名密码认证
    )
    )

    # 2. 创建一个信道(Channel),所有消息操作都在信道中进行
    ch = conn.channel()

    # 3. 开启生产者确认模式(Publisher Confirm)
    # 开启后,每条消息发送后都会等待 Broker 返回 ACK/NACK
    # 确保消息已被正确接收并持久化(如果设置了持久化)
    ch.confirm_delivery()

    # 4. 准备要发送的消息(Python 字典)
    message1 = {
    "id": 1,
    "task": "紧急备份",
    "priority": "high"
    }

    try:
    # 5. 发送消息到交换机
    # exchange: 交换机名称
    # routing_key: 路由键,交换机根据路由键将消息路由到对应的队列
    # body: 消息内容(JSON 字符串)
    # properties: 消息属性,这里设置 delivery_mode=2 表示消息持久化
    ch.basic_publish(
    exchange='my_exchange', # 交换机名称(需提前在 RabbitMQ 中声明或存在)
    routing_key='high', # 路由键(与队列绑定的 routing_key 对应)
    body=json.dumps(message1), # 将 Python 字典序列化为 JSON 字符串
    properties=pika.BasicProperties(
    delivery_mode=2 # 消息持久化(值为 2 表示持久化,1 表示非持久化)
    )
    )

    # 6. 如果没有抛出异常,说明 Broker 已确认接收消息
    print(f"✅ 发送高优先级消息成功: {message1} (已确认)")

    except pika.exceptions.UnroutableError:
    # 7. 捕获消息无法路由到任何队列的异常
    print(f"❌ 发送高优先级消息失败: {message1} (未路由到任何队列)")

    except Exception as e:
    # 8. 捕获其他可能的异常(如网络错误、认证失败等)
    print(f"❌ 发送高优先级消息失败: {e}")

    finally:
    # 9. 关闭连接
    conn.close()

    # ======================== 程序入口 ========================
    if __name__ == "__main__":
    send_message()

    3. consumer.py 

    #!/usr/bin/env python3
    # -*- coding: utf-8 -*-
    """
    consumer.py – RabbitMQ 消费者示例(交换机模式 + 手动确认)

    功能说明:
    1. 连接到 RabbitMQ Broker
    2. 声明队列(确保队列存在且持久化)
    3. 设置每次只处理一条消息(公平分发)
    4. 手动确认消息处理结果(ACK/NACK)
    5. 支持不同优先级任务的模拟处理

    作者:vb200811
    日期:2026-01-02
    """

    import pika
    import json
    import time

    # ======================== 配置参数 ========================
    HOST = '192.168.1.100' # RabbitMQ 服务器的 IP 地址
    PORT = 5672 # RabbitMQ 默认 AMQP 端口(5672)
    USER = 'mi_user' # 用户名
    PASSWORD = 'YourPassword123' # 密码
    QUEUE = 'high_queue' # 要监听的队列名称(可以改成 low_queue)

    # ======================== 消息处理回调函数 ========================
    def callback(ch, method, properties, body):
    """
    收到消息时调用的回调函数
    :param ch: 信道对象
    :param method: 消息方法属性(包含 delivery_tag 等)
    :param properties: 消息属性
    :param body: 消息内容(二进制数据)
    """
    try:
    # 1. 将二进制消息内容转换为 Python 字典
    message = json.loads(body.decode('utf-8'))
    print(f"\\n📨 收到消息: {message}")

    # 2. 模拟不同优先级任务的处理
    if message.get('priority') == 'high':
    print("🔧 处理高优先级任务…")
    time.sleep(2) # 模拟耗时操作
    else:
    print("🔧 处理普通任务…")
    time.sleep(1) # 模拟耗时操作

    print("✅ 处理完成")

    # 3. 手动确认消息(ACK)
    # 告诉 RabbitMQ 消息已成功处理,可以删除
    ch.basic_ack(delivery_tag=method.delivery_tag)

    except Exception as e:
    # 4. 处理失败时拒绝消息(NACK)
    # requeue=False 表示不重新入队,直接丢弃或进入死信队列
    print(f"❌ 处理失败: {e}")
    ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)

    # ======================== 启动消费者函数 ========================
    def start_consumer():
    """
    启动 RabbitMQ 消费者,监听指定队列
    """
    # 1. 创建连接(BlockingConnection 表示同步阻塞模式)
    conn = pika.BlockingConnection(
    pika.ConnectionParameters(
    host=HOST,
    port=PORT,
    credentials=pika.PlainCredentials(USER, PASSWORD) # 用户名密码认证
    )
    )

    # 2. 创建一个信道(Channel),所有消息操作都在信道中进行
    ch = conn.channel()

    # 3. 声明队列(确保队列存在)
    # durable=True 表示队列持久化,RabbitMQ 重启后队列仍然存在
    ch.queue_declare(queue=QUEUE, durable=True)

    # 4. 设置每次只处理一条消息(公平分发)
    # prefetch_count=1 表示消费者在处理完当前消息并确认之前,不再接收新的消息
    ch.basic_qos(prefetch_count=1)

    # 5. 开始消费消息
    # queue: 要监听的队列名称
    # on_message_callback: 收到消息时调用的回调函数
    # auto_ack=False: 关闭自动确认,开启手动确认模式
    ch.basic_consume(queue=QUEUE, on_message_callback=callback, auto_ack=False)

    print(f"🔄 开始监听队列: {QUEUE} ({HOST}:{PORT})")
    print("按 Ctrl+C 退出")

    try:
    # 6. 进入消费循环,持续监听队列
    ch.start_consuming()
    except KeyboardInterrupt:
    # 7. 捕获 Ctrl+C 中断信号,优雅退出
    print("\\n⏹️ 已停止")
    conn.close()

    # ======================== 程序入口 ========================
    if __name__ == "__main__":
    start_consumer()

    赞(0)
    未经允许不得转载:171主机测评 » Ubuntu 系统下 AMQP 协议 RabbitMQ服务器部署
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址