多线程向单个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
相关产品推荐
相关产品推荐

