Quarkus SmallRye RabbitMQ如何向默认交换机发送消息
SmallRye Reactive Messaging RabbitMQ 向默认交换机发消息实现方案
无需RabbitMQ服务端管控权限,通过消息元数据覆盖静态配置的方式即可实现需求,直接绕开MicroProfile Config无法配置空字符串交换机名的限制。
方案1:发送消息时附加元数据强制指定默认交换机(推荐)
SmallRye Reactive Messaging RabbitMQ 支持通过消息元数据动态覆盖配置文件中的交换机、路由键参数,元数据配置优先级高于静态配置,实现步骤如下:
- 保留现有静态channel的配置,不需要修改配置文件中的原有参数
- 发消息时构造
OutgoingRabbitMQMetadata对象,手动将exchange设置为空字符串,routing key设置为目标reply队列名称,附加到待发送消息上即可
代码示例:
import io.smallrye.reactive.messaging.rabbitmq.OutgoingRabbitMQMetadata; import org.eclipse.microprofile.reactive.messaging.Channel; import org.eclipse.microprofile.reactive.messaging.Emitter; import org.eclipse.microprofile.reactive.messaging.Message; import jakarta.inject.Inject; @Inject @Channel("your-rpc-reply-channel") Emitter<YourReplyPayload> replyEmitter; /** * 发送RPC响应到默认交换机 * @param payload 响应内容 * @param replyQueue 从请求消息中提取的reply-to队列名 * @param correlationId 从请求消息中提取的关联ID */ public void sendRpcResponse(YourReplyPayload payload, String replyQueue, String correlationId) { OutgoingRabbitMQMetadata rabbitMeta = OutgoingRabbitMQMetadata.builder() .withExchange("") // 强制指定为默认交换机,覆盖静态配置值 .withRoutingKey(replyQueue) // 路由键直接传入reply队列名 .withCorrelationId(correlationId) .build(); replyEmitter.send(Message.of(payload).addMetadata(rabbitMeta)); }
该方式完全符合RabbitMQ RPC模式规范,底层调用RabbitMQ客户端原生
basicPublish方法时会将exchange参数传空字符串,消息会被正确路由到指定名称的队列,不需要服务端做任何调整。
方案2:消费端自动回包场景适配
如果是通过@Incoming注解消费RPC请求、依靠方法返回值自动发送响应的场景,只需要在返回的Message对象中附加上述元数据即可,不要依赖框架自动生成的路由配置:
import org.eclipse.microprofile.reactive.messaging.Incoming; import org.eclipse.microprofile.reactive.messaging.Message; import io.smallrye.reactive.messaging.rabbitmq.IncomingRabbitMQMetadata; @Incoming("rpc-request-channel") public Message<YourReplyPayload> handleRpcRequest(Message<YourRequestPayload> requestMsg) { // 从入站消息中提取RPC标准参数 IncomingRabbitMQMetadata inMeta = requestMsg.getMetadata(IncomingRabbitMQMetadata.class) .orElseThrow(() -> new IllegalArgumentException("非标准RabbitMQ消息")); String replyQueue = inMeta.getReplyTo(); String correlationId = inMeta.getCorrelationId(); YourReplyPayload replyPayload = doBusinessProcess(requestMsg.getPayload()); // 构造回包元数据,指定走默认交换机 OutgoingRabbitMQMetadata outMeta = OutgoingRabbitMQMetadata.builder() .withExchange("") .withRoutingKey(replyQueue) .withCorrelationId(correlationId) .build(); return Message.of(replyPayload).addMetadata(outMeta); }
注意事项
- 不要尝试通过省略exchange配置的方式实现需求,框架未读取到exchange配置时,会默认使用当前channel的名称作为交换机名,消息会被发送到同名自定义交换机,无法路由到默认交换机绑定的队列
- 该方案不需要在RabbitMQ服务端新增交换机、绑定关系,完全适配无服务端管控权限的场景
- 兼容所有当前维护的Quarkus SmallRye RabbitMQ正式版本,无版本兼容问题
内容的提问来源于stack exchange,提问作者ReactiveMessenger
相关产品推荐
相关产品推荐

