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

如何实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 13:57:18