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

