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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 12:19:51