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

Solace JMS消费者重连后停止运行问题咨询

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 Connection and Session become invalid—they’re tied to the old broker instance.
  • Spring’s CachingConnectionFactory does handle connection caching and reconnection under the hood, but it won’t automatically refresh your manually created Session or MessageConsumer. 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 call commit() after successful processing and rollback() on failures to avoid message loss or duplicates.
  • Resource Cleanup: Always close stale Session and MessageConsumer instances 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 your CachingConnectionFactory setup.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 08:31:41