Solace中ConsumerFlowProperties为何延续生效?JmsListener重放全量消息
我正在审查一个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

