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

多线程向单个ActiveMQ Topic发消息时出现绑定异常的原因与解决方法

ActiveMQ批量发送消息出现Address already in use错误的原因与解决

问题背景

我们通过ActiveMQ实现两个Spring Boot服务间的通信:Service B向example-topic-1 Topic发送消息,Service A接收并处理。当处理约10000条IngestData数据时,采用每批1000条、10线程线程池的方式批量发送,运行一段时间后出现如下错误:

org.springframework.jms.UncategorizedJmsException: Uncategorized exception occurred during JMS processing; nested exception is javax.jms.JMSException: Could not connect to broker URL: tcp://localhost:61616. Reason: java.net.BindException: Address already in use: connect
    at org.springframework.jms.support.JmsUtils.convertJmsAccessException(JmsUtils.java:311) ~[spring-jms-5.3.23.jar:5.3.23]
    ...(省略中间栈信息)
Caused by: javax.jms.JMSException: Could not connect to broker URL: tcp://localhost:61616. Reason: java.net.BindException: Address already in use: connect
    at org.apache.activemq.util.JMSExceptionSupport.create(JMSExceptionSupport.java:38) ~[activemq-client-5.16.5.jar:5.3.23]

错误原因

  • 短连接耗尽本地端口:默认ActiveMQConnectionFactory创建的是短连接,每次调用jmsTemplate.convertAndSend()都会新建TCP连接,发送完成后关闭。但TCP连接关闭后会进入TIME_WAIT状态(默认2分钟),大量并发发送时,本地可用临时端口被快速耗尽,导致无法创建新连接。
  • 线程池放大连接压力:10个线程同时批量发送,每个线程内循环调用单条发送方法,短时间内创建大量TCP连接,加速端口耗尽。

解决方案

1. 配置连接池复用连接

使用ActiveMQ官方的PooledConnectionFactory替代原生ActiveMQConnectionFactory,通过连接池复用连接,避免频繁创建销毁连接:

修改ActiveMQConfiguration.java:

@Configuration
public class ActiveMQConfiguration {

    @Value("${activeMQ.broker-url}")
    private String brokerUrl;

    @Bean
    public ConnectionFactory connectionFactory(){
        ActiveMQConnectionFactory activeMQConnectionFactory = new ActiveMQConnectionFactory();
        activeMQConnectionFactory.setBrokerURL(brokerUrl);
        // 初始化连接池
        PooledConnectionFactory pooledConnectionFactory = new PooledConnectionFactory(activeMQConnectionFactory);
        pooledConnectionFactory.setMaxConnections(10); // 最大连接数匹配线程池大小
        pooledConnectionFactory.setIdleTimeout(30000); // 空闲连接超时时间(毫秒)
        return pooledConnectionFactory;
    }

    @Bean
    public JmsTemplate jmsTemplate(ActiveMQMessageConverter mActiveMQMessageConverter){
        JmsTemplate jmsTemplate = new JmsTemplate();
        jmsTemplate.setConnectionFactory(connectionFactory());
        jmsTemplate.setPubSubDomain(true);
        jmsTemplate.setMessageConverter(mActiveMQMessageConverter);
        return jmsTemplate;
    }
}

同时在build.gradle添加连接池依赖:

implementation 'org.apache.activemq:activemq-pool:5.16.5'

2. 优化批量发送逻辑

避免循环调用单条发送,改为复用Session批量发送,减少连接交互次数:

修改ActiveMQProducerService.java,新增批量发送方法:

@Component
@Slf4j
public class ActiveMQProducerService {

    @Autowired
    JmsTemplate jmsTemplate;

    @Value("${com.activeMQ.topic}")
    private String topic;

    // 单条发送(保留原方法,新增批量方法)
    public IngestData sendMessage(IngestData ingestData){
        try{
            jmsTemplate.convertAndSend(topic, ingestData);
        } catch(Exception e){
            log.error("Received Exception during sending for Message: "+ingestData.toString(), e);
        }
        return ingestData;
    }

    // 批量发送方法
    public void sendBatchMessage(List<IngestData> batchData){
        try{
            jmsTemplate.execute(session -> {
                MessageProducer producer = session.createProducer(new ActiveMQTopic(topic));
                for (IngestData data : batchData) {
                    Message message = session.createObjectMessage(data);
                    producer.send(message);
                }
                producer.close();
                return null;
            });
        } catch(Exception e){
            log.error("Received Exception during batch sending", e);
        }
    }
}

修改sendToActiveMQ方法,直接调用批量发送:

public void sendToActiveMQ(List<IngestData> batch){
    mPiActiveMQProducerService.sendBatchMessage(batch);
}

3. 调整系统TCP参数(可选)

如果端口耗尽问题仍存在,可调整系统TIME_WAIT相关参数,减少端口占用时间:

  • Linux系统:编辑/etc/sysctl.conf,添加以下配置后执行sysctl -p生效
    net.ipv4.tcp_tw_reuse = 1
    net.ipv4.tcp_tw_recycle = 1
    net.ipv4.tcp_fin_timeout = 30
    
  • Windows系统:修改注册表HKLM\SYSTEM\CurrentControlSet\Services\Tcpip\Parameters下的TcpTimedWaitDelay为30(十进制),重启系统生效

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 03:55:01