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
相关产品推荐
相关产品推荐

