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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:52:32