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

不依赖Spring Integration核心,使用TCP连接工厂处理自定义无分隔协议的问题

问题描述

我正在开发与某服务器的集成,该服务器基于原生TCP通信,采用无消息结束分隔符的自定义协议。我想了解是否可以不依赖Spring Integration其他组件,仅使用TcpNetClientConnectionFactory和CachingClientConnectionFactory。

我编写了如下测试代码:

public void process(Event event) {
    try {
        final TcpConnectionSupport connection = cachedClientCF.getConnection();
        final CountDownLatch countDownLatch = new CountDownLatch(1);
        connection.registerListener(message -> {
            final byte[] payload = (byte[]) message.getPayload();
            LOG.info("Payload: {}", new String(payload));
            countDownLatch.countDown();
            return true;
        });
        connection.send(new GenericMessage<>("Hello world"));
        countDownLatch.await();

    } catch (SoftEndOfStreamException e) {
        LOG.info("Stream closed because of EOF: " + e.getMessage());
    } catch (Exception e) {
        throw new RuntimeException(e);
    }
}

但遇到了问题:

  • 单次请求延迟可达5-7秒
  • 发送多个请求仅能收到一个响应
  • 监听器仅在连接关闭(收到FIN ACK)时才会触发,推测是因为没有配置消息结束标识

我想知道:若已知消息长度,是否可以保持连接打开,依次执行写、读、写、读操作?目前关于单独使用这些工厂的文档极少,查看源码后发现响应轮询速度极慢,每5秒才检查一次响应。

解决方案

完全可以仅通过这两个工厂实现长连接下的顺序读写,核心是配置固定长度消息解码器并调整轮询参数优化响应速度。

1. 配置固定长度消息解码器

既然已知消息长度,给TcpNetClientConnectionFactory设置固定长度解码器,框架就能精准识别消息边界,无需等待连接关闭才触发监听器:

TcpNetClientConnectionFactory clientFactory = new TcpNetClientConnectionFactory("serverHost", serverPort);
// 替换为你的实际固定消息长度
clientFactory.setCodec(TcpCodecs.fixedLength(100));
// 包装为缓存连接工厂实现连接复用
CachingClientConnectionFactory cachedClientCF = new CachingClientConnectionFactory(clientFactory);
// 根据并发需求设置连接池大小
cachedClientCF.setPoolSize(5);

2. 缩短响应轮询间隔

你遇到的5秒延迟是TcpNetConnection默认的soTimeout参数导致的,调整该参数可大幅降低轮询间隔:

// 设置为业务允许的最短超时时间,比如500毫秒
clientFactory.setSoTimeout(500);

这个参数控制Socket读取超时时间,同时会影响框架内部的响应检查频率,缩短后能有效减少等待延迟。

3. 长连接顺序读写注意事项

  • 配置固定长度解码器后,监听器会在收到完整的指定长度消息时立即触发,无需等待连接关闭
  • CachingClientConnectionFactory会自动复用连接,可在同一个连接上连续执行写、读操作
  • 多线程场景下,要么通过连接池分配独立连接,要么自己实现同步逻辑,避免读写混乱
  • 监听器返回true表示保持监听状态,刚好适配长连接下的多次读写需求

4. 调整后的示例代码

public void process(Event event) {
    try {
        final TcpConnectionSupport connection = cachedClientCF.getConnection();
        final CountDownLatch countDownLatch = new CountDownLatch(1);
        
        connection.registerListener(message -> {
            final byte[] payload = (byte[]) message.getPayload();
            LOG.info("Payload: {}", new String(payload));
            countDownLatch.countDown();
            return true; // 保持监听,支持后续请求
        });
        
        // 注意转成字节数组,匹配解码器的字节处理逻辑
        connection.send(new GenericMessage<>("Hello world".getBytes()));
        // 添加超时时间,避免无限等待
        countDownLatch.await(3, TimeUnit.SECONDS);

        // 可复用连接继续发送下一个请求
        // CountDownLatch secondLatch = new CountDownLatch(1);
        // connection.send(new GenericMessage<>("Second request".getBytes()));
        // secondLatch.await(3, TimeUnit.SECONDS);

    } catch (SoftEndOfStreamException e) {
        LOG.info("Stream closed because of EOF: " + e.getMessage());
        // 连接关闭后,缓存工厂会自动创建新连接供后续使用
    } catch (Exception e) {
        throw new RuntimeException(e);
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 07:19:59