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

请求:基于Apache ProtonJ2实现RabbitMQ Exchange的AMQP 1.0消费

Apache ProtonJ2 连接 RabbitMQ Exchange 消费(AMQP 1.0)实现示例

环境说明

  • 协议:AMQP 1.0
  • 客户端库:Apache ProtonJ2
  • Broker:RabbitMQ
  • 编程语言:Java 17

Maven 依赖

<dependency>
    <groupId>org.apache.qpid</groupId>
    <artifactId>protonj2-client</artifactId>
    <version>1.10.0</version> <!-- 使用最新稳定版即可 -->
</dependency>

实现代码

import org.apache.qpid.protonj2.client.*;
import org.apache.qpid.protonj2.client.exceptions.ClientException;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

public class ExchangeDirectConsumer {
    private static final String BROKER_HOST = "localhost";
    private static final int BROKER_PORT = 5672;
    private static final String TARGET_EXCHANGE = "your-target-exchange"; // 替换为实际Exchange名称
    private static final String USERNAME = "guest";
    private static final String PASSWORD = "guest";

    public static void main(String[] args) throws Exception {
        CountDownLatch runLatch = new CountDownLatch(1);

        // 初始化ProtonJ2客户端
        Client client = Client.create();

        try (Connection connection = client.connect(BROKER_HOST, BROKER_PORT, USERNAME, PASSWORD);
             Session session = connection.openSession()) {

            // 直接以Exchange名称作为Receiver目标地址
            Receiver receiver = session.openReceiver(TARGET_EXCHANGE, options -> {
                // 开启自动确认,也可根据业务需求改为手动确认
                options.autoAccept(true);
                // RabbitMQ会自动创建服务器命名的临时队列并绑定到目标Exchange
            });

            // 注册消息处理逻辑
            receiver.handler((delivery, message) -> {
                try {
                    String messageBody = message.body(String.class);
                    System.out.printf("Received message from Exchange '%s': %s%n", TARGET_EXCHANGE, messageBody);
                } catch (ClientException e) {
                    e.printStackTrace();
                }
                // 手动确认模式下需调用 delivery.accept()
            });

            // 启动Receiver开始消费
            receiver.open();

            // 保持程序运行,可根据业务场景调整退出逻辑
            runLatch.await(60, TimeUnit.SECONDS);

        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.err.println("Consumer process interrupted");
        } catch (ClientException e) {
            System.err.println("Consumer error occurred: " + e.getMessage());
            e.printStackTrace();
        } finally {
            client.close();
        }
    }
}

关键逻辑说明

  • 连接与会话:通过ProtonJ2的Client创建连接,自动管理Connection、Session资源,避免泄漏。
  • Receiver地址设置:直接将Receiver的目标地址指定为Exchange名称,RabbitMQ会自动完成:
    1. 创建服务器生成名称的临时队列(格式类似amq.gen-xxxxxx)
    2. 将临时队列绑定到目标Exchange
  • 消息确认:示例使用自动确认模式,若需精确控制消息消费状态,可改为手动确认并调用delivery.accept()。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 05:42:39