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

Spring Integration TCP服务器如何获取TcpConnection对象实现客户端消息回发?

Spring Integration TCP Server 获取TcpConnection对象问题

我用Spring Integration实现了一个TCP Socket Server,配置代码如下:

import java.util.concurrent.Executor;
import java.util.concurrent.Executors;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.config.EnableIntegration;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.dsl.IntegrationFlows;
import org.springframework.integration.ip.dsl.Tcp;
import org.springframework.integration.ip.dsl.TcpServerConnectionFactorySpec;
import org.springframework.integration.ip.tcp.connection.AbstractServerConnectionFactory;
import org.springframework.integration.ip.tcp.connection.TcpConnection;
import org.springframework.integration.ip.tcp.serializer.ByteArrayLfSerializer;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.support.ErrorMessage;


@Configuration
@EnableIntegration
public class TcpServerConfig
{

  private static final Logger logger = LoggerFactory.getLogger(TcpServerConfig.class);
  public static final String IP_HEADER = "ip_connection";

  public TcpServerConfig()
  {}

  @Bean
  public AbstractServerConnectionFactory serverFactory()
  {
    TcpServerConnectionFactorySpec factory = Tcp.netServer(4444)
      .serializer(new ByteArrayLfSerializer())
      .deserializer(new ByteArrayLfSerializer());
    return factory.get();
  }

  @Bean
  public MessageChannel inboundChannel()
  {
    return new DirectChannel();
  }

  @Bean
  public IntegrationFlow inboundFlow()
  {
    return IntegrationFlows.from(Tcp.inboundAdapter(serverFactory()))
      .channel(inboundChannel())
      .get();
  }

  @Bean
  public Executor executor()
  {
    return Executors.newCachedThreadPool();
  }

  @Bean
  @ServiceActivator(inputChannel = "inboundChannel")
  public MessageHandler inboundHandler(TcpSocketServer aTcpSocketServerHandler)
  {
    return message -> {
      try
      {
        byte[] payload = (byte[]) message.getPayload();
        String data = new String(payload);

        MessageHeaders headers = message.getHeaders();
        headers.forEach((key, value) -> System.out.println(key + ": " + value));

        TcpConnection connection = headers.get("ip_connectionId", TcpConnection.class);

        logger.info("Received message: {}", data);
        aTcpSocketServerHandler.handleMessage(data, connection);
      }
      catch (Exception e)
      {
        logger.error("Error processing message", e);
        throw e;
      }
    };
  }

  @Bean
  @ServiceActivator(inputChannel = "errorChannel")
  public MessageHandler errorHandler()
  {
    return message -> {
      ErrorMessage errorMessage = (ErrorMessage) message;
      logger.error("Error occurred: ", errorMessage.getPayload());
    };
  }
}

测试用的TCP客户端代码:

public static void main(String[] args)
  {
    String hostname = "localhost";
    int port = 4444;
    String message = "Hello, TCP Server!";

    try (Socket socket = new Socket(hostname, port))
    {

      // Send message to the server
      OutputStream output = socket.getOutputStream();
      output.write((message + "\n").getBytes());
      output.flush();

      // Read response from the server
      BufferedReader reader = new BufferedReader(new InputStreamReader(socket.getInputStream()));
      String response;
      while ((response = reader.readLine()) != null)
      {
        System.out.println("Server response: " + response);
      }

    }
    catch (IOException ex)
    {
      System.err.println("I/O error: " + ex.getMessage());
    }
  }

查看MessageHeaders得到以下结果:

  • ip_connectionId: 127.0.0.1:53129:4444:e9aad42f-a057-4c02-87fc-0cc9a60427fe
  • ip_localInetAddress: /127.0.0.1
  • ip_address: 127.0.0.1
  • id: 1ad8f0c5-a882-83d6-8f55-ceddce8961bf
  • ip_hostname: 127.0.0.1
  • timestamp: 1718618786314

可以看到ip_connectionId是字符串类型,无法直接获取TcpConnection对象。我需要把这个对象注入到TcpSocketServer处理器中实现向客户端回发消息,请问用这种Spring Integration的实现方式能不能获取到TcpConnection对象?


解决方案

当然可以获取到TcpConnection对象,以下两种方式都能实现:

方式一:通过连接工厂根据connectionId获取

注入你的AbstractServerConnectionFactory实例,调用它的getConnection(String connectionId)方法,传入消息头里的ip_connectionId字符串即可拿到对应的TcpConnection对象:

修改inboundHandler代码:

@Bean
@ServiceActivator(inputChannel = "inboundChannel")
public MessageHandler inboundHandler(TcpSocketServer aTcpSocketServerHandler, AbstractServerConnectionFactory serverFactory) {
    return message -> {
        try {
            byte[] payload = (byte[]) message.getPayload();
            String data = new String(payload);

            MessageHeaders headers = message.getHeaders();
            String connectionId = headers.get("ip_connectionId", String.class);
            // 通过连接工厂获取TcpConnection对象
            TcpConnection connection = serverFactory.getConnection(connectionId);

            logger.info("Received message: {}", data);
            aTcpSocketServerHandler.handleMessage(data, connection);
        } catch (Exception e) {
            logger.error("Error processing message", e);
            throw e;
        }
    };
}

方式二:关闭负载提取,直接获取TcpMessage

默认情况下,TCP入站适配器会提取消息负载,仅将连接ID放入消息头。设置extractPayload=false后,消息负载会是一个TcpMessage对象,该对象直接包含TcpConnection和原始数据:

修改inboundFlow配置:

@Bean
public IntegrationFlow inboundFlow() {
    TcpReceivingChannelAdapter adapter = Tcp.inboundAdapter(serverFactory())
            .extractPayload(false) // 关闭负载提取,保留完整TcpMessage
            .get();
    return IntegrationFlows.from(adapter)
            .channel(inboundChannel())
            .get();
}

然后修改inboundHandler处理TcpMessage:

@Bean
@ServiceActivator(inputChannel = "inboundChannel")
public MessageHandler inboundHandler(TcpSocketServer aTcpSocketServerHandler) {
    return message -> {
        try {
            TcpMessage tcpMessage = (TcpMessage) message.getPayload();
            byte[] payloadBytes = tcpMessage.getPayload();
            String data = new String(payloadBytes);
            // 直接从TcpMessage获取TcpConnection
            TcpConnection connection = tcpMessage.getConnection();

            logger.info("Received message: {}", data);
            aTcpSocketServerHandler.handleMessage(data, connection);
        } catch (Exception e) {
            logger.error("Error processing message", e);
            throw e;
        }
    };
}

拿到TcpConnection对象后,调用它的send(Message<?> message)方法就能向客户端回发消息。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 07:48:10