Spring Integration:实现队列消息的选择性消费者分发
咱们来搞定这个Spring Integration里选择性消费者的消息留存问题!先梳理下你的场景和现有问题,再一步步给出调整方案。
问题场景与现有实现
你需要实现一个Spring Integration流,让消息留存到匹配的消费者出现后再被消费,目前的实现分为三个部分:
1. 消息接收与入队Flow
@Bean public IntegrationFlow messagesPerCustomerFlow() { return IntegrationFlows. from(WebFlux.inboundChannelAdapter("/messages/{customer}") .requestMapping(r -> r .methods(HttpMethod.POST) ) .requestPayloadType(JsonNode.class) .headerExpression("customer", "#pathVariables.customer") ) .channel(messagesPerCustomerQueue()) .get(); } @Bean(name = PollerMetadata.DEFAULT_POLLER) public PollerSpec poller() { return Pollers.fixedRate(100); } @Bean public QueueChannel messagesPerCustomerQueue() { return MessageChannels.queue() .get(); }
2. SSE推送Flow
@Bean public IntegrationFlow eventMessagesPerCustomer() { return IntegrationFlows .from(WebFlux.inboundGateway("/events/{customer}") .requestMapping(m -> m.produces(TEXT_EVENT_STREAM_VALUE)) .headerExpression("customer", "#pathVariables.customer") .payloadExpression("''") // 为了让handle((p,h)正常工作 ) .log() .handle((p, h) -> { String customer = h.get("customer").toString(); PublisherSubscription<JsonNode> publisherSubscription = subscribeToMessagesPerCustomer(customer); return Flux.from(publisherSubscription.getPublisher()) .map(Message::getPayload) .doFinally(signalType -> publisherSubscription.unsubscribe()); }) .get(); }
3. 动态注册消费者Flow
public PublisherSubscription<JsonNode> subscribeToMessagesPerCustomer(String customer) { IntegrationFlowBuilder flow = IntegrationFlows.from(messagesPerCustomerQueue()) .filter("headers.customer=='" + customer + "'", filterEndpointSpec -> filterEndpointSpec.throwExceptionOnRejection(true)); Publisher<Message<JsonNode>> messagePublisher = flow.toReactivePublisher(); IntegrationFlowRegistration registration = integrationFlowContext.registration(flow.get()) .register(); return new PublisherSubscription<>(messagePublisher, registration); }
当前遇到的问题
- 无任何订阅者时发送消息,抛出
MessageDeliveryException: Dispatcher has no subscribers for channel 'application.messagesPerCustomerQueue' - 无匹配订阅者时发送消息,抛出
AggregateMessageDeliveryException: All attempts to deliver Message to MessageHandlers failed
解决方案
要实现消息留存到匹配消费者出现或过期的目标,需要调整两个核心配置:
1. 修改队列通道,允许无订阅者时留存消息
默认的QueueChannel在没有订阅者时会直接拒绝消息并抛出异常,我们需要调整它的Dispatcher配置,让消息能存入队列而不是报错:
@Bean public QueueChannel messagesPerCustomerQueue() { QueueChannel queueChannel = new QueueChannel(); // 获取队列的Dispatcher,关闭无订阅者时的失败抛出 UnicastingDispatcher dispatcher = (UnicastingDispatcher) queueChannel.getDispatcher(); dispatcher.setFailOnNoSubscribers(false); return queueChannel; }
2. 调整Filter的异常策略,让不匹配的消息回到队列
当前设置的throwExceptionOnRejection(true)会导致不匹配的消息直接抛出异常,进而触发全局投递失败。我们需要关闭这个配置,让不匹配的消息自动回到队列等待下一次轮询:
public PublisherSubscription<JsonNode> subscribeToMessagesPerCustomer(String customer) { IntegrationFlowBuilder flow = IntegrationFlows.from(messagesPerCustomerQueue()) // 关闭拒绝时的异常抛出,不匹配的消息会回到队列 .filter("headers.customer=='" + customer + "'", spec -> spec.throwExceptionOnRejection(false)) .log("匹配到客户[" + customer + "]的消息"); Publisher<Message<JsonNode>> messagePublisher = flow.toReactivePublisher(); IntegrationFlowRegistration registration = integrationFlowContext.registration(flow.get()) .register(); return new PublisherSubscription<>(messagePublisher, registration); }
3. 可选:添加消息过期机制
如果需要让消息在队列中留存一段时间后自动过期,可以在消息入队时添加过期头:
@Bean public IntegrationFlow messagesPerCustomerFlow() { return IntegrationFlows. from(WebFlux.inboundChannelAdapter("/messages/{customer}") .requestMapping(r -> r.methods(HttpMethod.POST)) .requestPayloadType(JsonNode.class) .headerExpression("customer", "#pathVariables.customer") ) // 设置消息1小时后过期(单位:毫秒) .transform((Message<JsonNode> msg) -> MessageBuilder.fromMessage(msg) .setHeader(MessageHeaders.EXPIRATION_DATE, System.currentTimeMillis() + 3600000) .build()) .channel(messagesPerCustomerQueue()) .get(); }
之后可以配合定时任务,定期清理队列中的过期消息,或者使用Spring Integration的MessageStore来管理带过期策略的消息。
原理说明
setFailOnNoSubscribers(false):让队列通道在没有订阅者时,将消息存入队列而非抛出异常,实现基础的消息留存。- 关闭
throwExceptionOnRejection:不匹配的消息会自动回到队列,等待下一次轮询,直到有匹配的消费者接收。 - 过期头配置:给消息设置存活时间,避免无效消息长期占用队列资源。
调整后就能实现标准的选择性消费者模式:消息留存于队列中,反复推送给订阅者直至被消费或过期。
内容的提问来源于stack exchange,提问作者dschulten
相关产品推荐
相关产品推荐

