Solace JMS消费者重连后停止运行问题咨询
Alright, let's tackle this problem head-on. I’ve run into similar issues with Solace JMS and Spring’s CachingConnectionFactory before, so I know exactly what’s going on here.
问题根源
Your current setup manually creates a Connection, Session, and MessageConsumer, then hands the consumer off to a separate thread. Here’s why this breaks after failover:
- When Solace fails over, the original
ConnectionandSessionbecome invalid—they’re tied to the old broker instance. - Spring’s
CachingConnectionFactorydoes handle connection caching and reconnection under the hood, but it won’t automatically refresh your manually createdSessionorMessageConsumer. Your thread keeps trying to use the stale resources, which leads to silent failures or exceptions that stop consumption.
解决方案1:让Spring完全管理消费者(推荐)
The easiest and most reliable fix is to let Spring’s JMS infrastructure manage your consumers, sessions, and connections. This way, Spring automatically handles failover, reconnections, and resource cleanup for you.
步骤1:配置JMS监听器容器工厂
First, create a JmsListenerContainerFactory tailored to your Solace setup (with transaction support):
@Bean public JmsListenerContainerFactory<?> solaceJmsListenerContainerFactory(CachingConnectionFactory ccf) { DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); factory.setConnectionFactory(ccf); factory.setSessionTransacted(true); // Matches your SESSION_TRANSACTED setting factory.setConcurrency("1"); // Ensures a single consumer thread (matches your original setup) factory.setAutoStartup(true); return factory; }
步骤2:使用@JmsListener定义消费者
Replace your manual consumer creation with a Spring-managed listener. Spring will handle failover and reconnection automatically:
@Component public class SolaceMessageHandler { @JmsListener( destination = "your-target-destination", containerFactory = "solaceJmsListenerContainerFactory" ) public void processMessage(Message message, Session session) throws JMSException { try { // Your message processing logic goes here System.out.println("Received message: " + ((TextMessage) message).getText()); // Commit the transaction if processing succeeds session.commit(); } catch (Exception e) { // Rollback on failure to redeliver the message session.rollback(); // Add error logging or alerting here e.printStackTrace(); } } }
解决方案2:手动管理资源(如果必须保留自定义线程)
If you absolutely need to manage the consumer thread yourself, you’ll need to add connection event listening to detect failover and refresh your Session/Consumer when a reconnection happens.
实现代码示例
@Component public class ManualSolaceConsumer { @Autowired private CachingConnectionFactory ccf; private volatile MessageConsumer consumer; private volatile Session session; private volatile Connection connection; private Thread consumerThread; @PostConstruct public void initConsumer() { setupConsumerResources(); } private void setupConsumerResources() { try { // Get a connection from the caching factory connection = ccf.createConnection(); // Add Solace-specific connection listener to detect reconnections if (connection instanceof com.solacesystems.jms.SolConnection) { ((com.solacesystems.jms.SolConnection) connection).addConnectionListener(new com.solacesystems.jms.SolConnectionListener() { @Override public void onConnectionEvent(com.solacesystems.jms.SolConnectionEvent event) { if (event.getEventType() == com.solacesystems.jms.SolConnectionEvent.EventType.RECONNECTED) { // Reconnection successful—reset consumer resources resetConsumer(); } } }); } connection.start(); // Create transacted session (matches your original setup) session = connection.createSession(true, Session.SESSION_TRANSACTED); Destination destination = session.createQueue("your-target-destination"); consumer = session.createConsumer(destination); // Start the consumer thread startConsumerThread(); } catch (JMSException e) { e.printStackTrace(); // Add retry logic here if needed } } private void startConsumerThread() { consumerThread = new Thread(() -> { while (!Thread.currentThread().isInterrupted()) { try { Message message = consumer.receive(); if (message != null) { // Process your message System.out.println("Manual thread received message: " + ((TextMessage) message).getText()); session.commit(); } } catch (JMSException e) { // Check if the exception is connection-related (failover) if (isConnectionFailure(e)) { break; // Exit loop to trigger reset } // Rollback on processing errors try { session.rollback(); } catch (JMSException ex) { ex.printStackTrace(); } } } }); consumerThread.start(); } private void resetConsumer() { // Clean up stale resources try { if (consumer != null) consumer.close(); if (session != null) session.close(); if (consumerThread != null) consumerThread.interrupt(); } catch (JMSException e) { e.printStackTrace(); } // Recreate consumer resources setupConsumerResources(); } private boolean isConnectionFailure(JMSException e) { // Identify Solace-specific connection error codes String errorCode = e.getErrorCode(); return errorCode != null && (errorCode.contains("CONNECTION_LOST") || errorCode.contains("RECONNECT_ATTEMPT")); } }
关键注意事项
- Transaction Handling: When using
SESSION_TRANSACTED, always ensure you callcommit()after successful processing androllback()on failures to avoid message loss or duplicates. - Resource Cleanup: Always close stale
SessionandMessageConsumerinstances before recreating them to prevent resource leaks. - Solace Connection Properties: Make sure your Solace JMS client is configured with proper reconnect settings (e.g.,
reconnectRetries,reconnectRetryWaitInMillis) in yourCachingConnectionFactorysetup.
内容的提问来源于stack exchange,提问作者Daniel Urbaniak

