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

Solace中ConsumerFlowProperties为何延续生效?JmsListener重放全量消息

问题:Solace队列重放设置为何会影响后续JMS消费者?

我正在审查一个Solace应用的代码库:先通过JCSMP配置FlowReceiver并设置ConsumerFlowProperties,指定从队列起始位置重放消息,启动后直接关闭FlowReceiver且未读取任何消息,随后关闭会话;最后启动关联@JmsListener注解方法的ListenerContainer时,却发现它仍从队列起始位置重放所有消息。为何ConsumerFlowProperties的重放设置会延续生效?它不应仅关联之前的会话/消费者吗?

代码示例:

public void replayThenStart() {
    final var jcsmpFactory = JCSMPFactory.onlyInstance();
    JCSMPSession session = null;
    try {
        session = jcsmpFactory.createSession(jcsmpProperties);
        session.connect();

        final var queue = jcsmpFactory.createQueue("q/uat/event");
        final var consumerFlowProperties = new ConsumerFlowProperties();
        consumerFlowProperties.setEndpoint(queue);
        consumerFlowProperties.setReplayStartLocation(jcsmpFactory.createReplayStartLocationBeginning());
        consumerFlowProperties.setActiveFlowIndication(false);

        FlowReceiver consumer = null;
        try {
            consumer = session.createFlow(this, consumerFlowProperties);
            consumer.start();
            log.info("Replay flow (" + consumer + ") created.");
        } catch (Throwable t) {
            log.error("Failed replay flow.", t);
        } finally {
            if (consumer != null) {
                log.info("Close replay flow.");
                consumer.close();
            }
        }
    } catch (Throwable t) {
        log.error("Failed replay session.", t);
    } finally {
        if (session != null) {
            session.closeSession();
        }
    }

    Objects.requireNonNull(jmsListenerEndpointRegistry.getListenerContainer("eventListener")).start();
}

@JmsListener(id = "eventListener", destination = "q/uat/event", containerFactory = "eventContainerFactory")
public void onEvents(final List<Event> events, @Header("authentication") final String header) {
    try {
        onEvents(events, jwtUtil.getAuthentication(header));
    } catch (Throwable e) {
        log.error("Failed authenticating header [{}].", header, e);
        log.debug("Unable to retrieve events - [{}]", events);
    }
}

public void onEvents (final List<Event> events, final Authentication authentication) {
    //process events
}

原因分析

问题的核心在于:Solace队列的重放起点设置是队列级别的全局属性,并非仅绑定到单个消费者或会话。

当你通过JCSMP的setReplayStartLocation(Beginning)创建并启动FlowReceiver时,这个操作会直接修改目标队列的重放配置——队列被标记为需要从起始位置开始重放。即使你立刻关闭了FlowReceiver和会话,队列的这个重放状态并不会自动重置。后续任何连接到该队列的消费者(包括你的JMS Listener)都会继承这个队列级别的重放设置,从而从队列开头开始消费消息。

另外,setActiveFlowIndication(false)只是让FlowReceiver处于非活跃状态(不主动拉取消息),但它并不会阻止队列重放配置的修改——只要Flow被创建并启动,队列的重放起点就已经被更新了。

解决方案

根据你的需求,有两种可行的解决方式:

1. 临时设置重放后,手动重置队列配置

在关闭第一个FlowReceiver后,创建一个新的FlowReceiver,将队列的重放起点重置为当前位置,这样后续消费者会从正常位置开始消费:

// 在关闭第一个consumer之后添加以下代码
try {
    // 创建重置用的FlowProperties
    ConsumerFlowProperties resetProps = new ConsumerFlowProperties();
    resetProps.setEndpoint(queue);
    resetProps.setReplayStartLocation(jcsmpFactory.createReplayStartLocationCurrent());
    resetProps.setActiveFlowIndication(false);
    
    FlowReceiver resetConsumer = session.createFlow(this, resetProps);
    resetConsumer.start();
    log.info("Reset queue replay location to current.");
    resetConsumer.close();
} catch (Exception e) {
    log.error("Failed to reset queue replay location.", e);
}

2. 直接在目标消费者(JMS Listener)中配置重放

如果你本来就希望JMS Listener从队列起始位置重放,那完全不需要先用JCSMP做临时操作,直接在JMS容器工厂中配置Solace的重放扩展属性即可:

// 配置JMS容器工厂时添加重放属性
@Bean
public JmsListenerContainerFactory<?> eventContainerFactory(ConnectionFactory connectionFactory) {
    DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory();
    factory.setConnectionFactory(connectionFactory);
    
    // 设置Solace JMS扩展属性,指定从队列开头重放
    factory.setSessionTransacted(false);
    factory.setSessionAcknowledgeMode(Session.CLIENT_ACKNOWLEDGE);
    
    // 获取Solace专属的ConnectionFactory,设置重放属性
    if (connectionFactory instanceof SolConnectionFactory) {
        ((SolConnectionFactory) connectionFactory).setReplayStartLocation(ReplayStartLocation.BEGINNING);
    }
    
    return factory;
}

这样启动JMS Listener时,会直接应用重放配置,无需依赖临时的JCSMP会话。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 09:15:03