Java Socket编程遇java.net.ConnectException超时问题求助
问题描述
我是Java新手,正在构建一个微服务,该服务从ActiveMQ读取消息并通过Socket编程将消息作为TCP数据包发送至另一台服务器。目前能成功从ActiveMQ读取消息,但发送少量TCP数据包后出现java.net.ConnectException: Operation timed out异常,求解决办法。
相关代码
App.java
package com.example; import java.io.File; import java.io.FileWriter; import java.io.BufferedWriter; import java.io.IOException; import javax.jms.*; import java.io.*; import java.net.Socket; import org.apache.activemq.ActiveMQConnectionFactory; public class App { private static String url = "tcp://172.13.0.1:61616"; private static String topicName = "topic/CNAAlarms"; public static void main(String[] args) { try { connectToActiveMQ(url, topicName); }catch(Exception e){ System.out.println("for connectToActiveMQ == " + e.getMessage()); System.out.println("------------"); } } public static void connectToActiveMQ(String url, String topicName) throws JMSException { System.out.println("inside connectToActiveMQ"); System.out.println("url = " + url); System.out.println("topicName = " + topicName); ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(); connectionFactory.setUserName("abc@abc"); connectionFactory.setPassword("abc"); connectionFactory.setClientID("testingClient"); connectionFactory.setBrokerURL(url); Connection connection = connectionFactory.createConnection(); System.out.println("connected"); connection.start(); Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); System.out.println("session created"); Topic topic = session.createTopic(topicName); TopicSubscriber durablSubscriber = session.createDurableSubscriber(topic, topicName); System.out.println("durable subscriber created"); MessageListener listener = new Listener(); durablSubscriber.setMessageListener(listener); } public static void sendDataToTCPServer(String data){ // initialize socket and output streams try { System.out.println("Inside SendDataToTCPServer"); Socket s = new Socket("abc123d.it.internal", 11201); DataOutputStream dos = new DataOutputStream(s.getOutputStream()); dos.writeUTF(data); dos.close(); s.close(); }catch(IOException e) { e.printStackTrace(); } } }
Listener.java
package com.example; import javax.jms.*; public class Listener implements MessageListener{ public void onMessage(Message message){ try{ if (message instanceof TextMessage) { TextMessage textMessage = (TextMessage) message; String text = textMessage.getText(); System.out.println("Received text: " + text); App.sendDataToTCPServer("activemq-bridge :" + text + "\n"); } else { System.out.println("Received: " + message); App.sendDataToTCPServer("activemq-bridge :" +message + "\n"); } }catch(Exception e){ System.out.println("for onMessage Listener: " + e.getMessage()); } } }
报错信息
Received text: {"userName":null,"enterpriseName":"ABC"} Inside SendDataToTCPServer java.net.ConnectException: Operation timed out at java.base/sun.nio.ch.Net.connect0(Native Method) at java.base/sun.nio.ch.Net.connect(Net.java:579) at java.base/sun.nio.ch.Net.connect(Net.java:568) at java.base/sun.nio.ch.NioSocketImpl.connect(NioSocketImpl.java:585) at java.base/java.net.SocksSocketImpl.connect(SocksSocketImpl.java:327) at java.base/java.net.Socket.connect(Socket.java:633) at java.base/java.net.Socket.connect(Socket.java:583) at java.base/java.net.Socket.<init>(Socket.java:507) at java.base/java.net.Socket.<init>(Socket.java:287) at com.example.App.sendDataToTCPServer(App.java:98) at com.example.Listener.onMessage(Listener.java:12) at org.apache.activemq.ActiveMQMessageConsumer.dispatch(ActiveMQMessageConsumer.java:1321) at org.apache.activemq.ActiveMQSessionExecutor.dispatch(ActiveMQSessionExecutor.java:131) at org.apache.activemq.ActiveMQSessionExecutor.iterate(ActiveMQSessionExecutor.java:202) at org.apache.activemq.thread.PooledTaskRunner.runTask(PooledTaskRunner.java:129) at org.apache.activemq.thread.PooledTaskRunner$1.run(PooledTaskRunner.java:47) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:833)
解决办法
核心问题分析
当前代码的主要问题是每处理一条消息就新建并销毁Socket连接,短时间内大量请求会导致:
- TCP连接资源耗尽,系统无法创建新连接
- 目标服务器触发连接频率限制,拒绝新连接
- 无连接超时配置,等待时间过长引发超时
具体修复方案
复用Socket长连接
将Socket连接改为单例复用,避免频繁创建销毁:
package com.example; import java.io.*; import java.net.InetSocketAddress; import java.net.Socket; import javax.jms.*; import org.apache.activemq.ActiveMQConnectionFactory; public class App { private static String url = "tcp://172.13.0.1:61616"; private static String topicName = "topic/CNAAlarms"; // TCP连接相关配置 private static Socket tcpSocket; private static DataOutputStream dos; private static final String TCP_HOST = "abc123d.it.internal"; private static final int TCP_PORT = 11201; private static final int CONNECT_TIMEOUT = 5000; // 5秒连接超时 public static void main(String[] args) { try { connectToActiveMQ(url, topicName); // 程序退出时关闭TCP连接 Runtime.getRuntime().addShutdownHook(new Thread(App::closeTcpConnection)); } catch(Exception e) { System.out.println("for connectToActiveMQ == " + e.getMessage()); System.out.println("------------"); } } public static void connectToActiveMQ(String url, String topicName) throws JMSException { System.out.println("inside connectToActiveMQ"); System.out.println("url = " + url); System.out.println("topicName = " + topicName); ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(); connectionFactory.setUserName("abc@abc"); connectionFactory.setPassword("abc"); connectionFactory.setClientID("testingClient"); connectionFactory.setBrokerURL(url); Connection connection = connectionFactory.createConnection(); System.out.println("connected"); connection.start(); Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); System.out.println("session created"); Topic topic = session.createTopic(topicName); TopicSubscriber durablSubscriber = session.createDurableSubscriber(topic, topicName); System.out.println("durable subscriber created"); MessageListener listener = new Listener(); durablSubscriber.setMessageListener(listener); } // 初始化或复用TCP连接 private static void initTcpConnection() throws IOException { if (tcpSocket == null || !tcpSocket.isConnected() || tcpSocket.isClosed()) { tcpSocket = new Socket(); tcpSocket.connect(new InetSocketAddress(TCP_HOST, TCP_PORT), CONNECT_TIMEOUT); dos = new DataOutputStream(tcpSocket.getOutputStream()); } } public static void sendDataToTCPServer(String data){ try { System.out.println("Inside SendDataToTCPServer"); initTcpConnection(); dos.writeUTF(data); dos.flush(); // 确保数据立即发送 } catch(IOException e) { e.printStackTrace(); // 连接异常时重置资源 resetTcpConnection(); } } // 重置TCP连接资源 private static void resetTcpConnection() { try { if (dos != null) dos.close(); if (tcpSocket != null) tcpSocket.close(); } catch (IOException e) { e.printStackTrace(); } tcpSocket = null; dos = null; } // 关闭TCP连接 private static void closeTcpConnection() { resetTcpConnection(); } }
添加连接重试逻辑(可选)
如果目标服务器偶尔不稳定,可以添加重试机制:
private static void initTcpConnection() throws IOException { int retryCount = 3; IOException lastError = null; while (retryCount > 0) { try { if (tcpSocket == null || !tcpSocket.isConnected() || tcpSocket.isClosed()) { tcpSocket = new Socket(); tcpSocket.connect(new InetSocketAddress(TCP_HOST, TCP_PORT), CONNECT_TIMEOUT); dos = new DataOutputStream(tcpSocket.getOutputStream()); return; } } catch (IOException e) { lastError = e; retryCount--; try { Thread.sleep(1000); // 重试间隔1秒 } catch (InterruptedException ie) { Thread.currentThread().interrupt(); } } } throw lastError; }
检查目标服务器配置
- 确认目标服务器
abc123d.it.internal的11201端口对外开放,防火墙允许当前服务器访问 - 检查目标服务器的TCP连接数限制,是否因短时间内大量连接被拒绝
- 确认目标服务器运行正常,未因负载过高无法响应
优化消息消费线程(可选)
调整ActiveMQ消费线程池大小,避免线程资源耗尽:
ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(); // 设置消费线程池参数 connectionFactory.setThreadPoolMaxSize(10); connectionFactory.setThreadPoolMinSize(2);
内容的提问来源于stack exchange,提问作者user61815
相关产品推荐
相关产品推荐

