如何实现Spring Integration Flow作为TCP客户端并接收消息?
Spring Integration TCP客户端配置问题
我尝试配置Spring Integration Flow作为TCP客户端,连接运行在localhost:9026的服务器。此前使用简单Socket实现该功能成功,但Spring Integration的配置却无法正常运行。我期望应用启动时自动连接服务器,接收心跳或XML格式的消息,目前无法实现,恳请提供帮助。
初始代码
@SpringBootApplication @EnableIntegration @EnableIntegrationManagement(defaultLoggingEnabled = "true") public class TcplistenerApplication { private static final Logger LOGGER = LoggerFactory.getLogger(TcplistenerApplication.class); public static void main(String[] args) throws IOException { SpringApplication.run(TcplistenerApplication .class, args); @Bean public AbstractConnectionFactory clientConnectionFactory() { TcpNetClientConnectionFactory factory = new TcpNetClientConnectionFactory("localhost", 9026); factory.setDeserializer(TcpCodecs.raw()); return factory; } @Bean public TcpSendingMessageHandler inbound() { TcpSendingMessageHandler adapter = new TcpSendingMessageHandler(); adapter.setConnectionFactory(clientConnectionFactory()); adapter.setClientMode(true); return adapter; } @Bean public IntegrationFlow tcpClientFlow() { return IntegrationFlows.from(Tcp.inboundAdapter(clientConnectionFactory())) .handle(m -> LOGGER.info(m.getPayload().toString())).get(); } @EventListener public void listen(TcpConnectionEvent event) { LOGGER.info(event.toString()); } }
更新后的配置(已可接收数据,第一步仅实现日志打印)
//... @Bean public AbstractConnectionFactory clientConnectionFactory() { TcpNetClientConnectionFactory factory = new TcpNetClientConnectionFactory("localhost", 9026); factory.setDeserializer(new CustomMessageDeserializer()); return factory; } @Bean public IntegrationFlow tcpClientFlow() { return IntegrationFlows.from(Tcp.inboundAdapter(clientConnectionFactory()).clientMode(true)).log() .handle(m -> Logger.info(m.getPayload())).get(); } //...
CustomMessageDeserializer代码
public class CustomMessageDeserializer implements Deserializer<CustomMessage> { private CustomMessageParser parser; public CustomMessageDeserializer() { this.parser = new CustomMessageParser(); } @Override public CustomMessage deserialize(InputStream inputStream) throws IOException { byte[] buffer = new byte[1024]; int readBytes = inputStream.read(buffer); return parser.parse(Arrays.copyOf(buffer, readBytes)); } }
内容的提问来源于stack exchange,提问作者Galadriel
相关产品推荐
相关产品推荐

