如何用Apache Camel实现带持久连接、异步响应匹配的Netty TCP客户端?
基于Apache Camel实现异步TCP请求响应匹配的持久连接客户端方案
核心实现逻辑
要实现异步响应与请求的匹配,核心是基于唯一请求ID关联机制,配合Camel的聚合EIP做上下文暂存与匹配,同时配置Netty组件开启长连接复用,具体步骤如下:
具体实现步骤
第一步:Netty端点配置开启持久长连接
修改原有Netty端点参数,禁用短连接、开启长连接保活,配置异步模式:
// 参数说明:sync=false开启异步模式,keepAlive保持长连接复用,tcpNoDelay降低延迟 netty:tcp://127.0.0.1:9898?sync=false&keepAlive=true&tcpNoDelay=true&reuseAddress=true&clientMode=true
第二步:请求添加唯一关联ID
和TCP服务端约定:所有请求、响应报文中必须携带同一个唯一请求ID,在报文转换处理器中生成并存储ID:
private void transformTcpMessage(Exchange exchange) { RequestBean request = exchange.getIn().getBody(RequestBean.class); // 生成全局唯一请求ID String requestId = UUID.randomUUID().toString(); // 存入Exchange Header留作后续匹配使用 exchange.getIn().setHeader("TCP_REQUEST_ID", requestId); // 将requestId写入发送给服务端的TCP报文中 TcpReqPacket packet = new TcpReqPacket(); packet.setRequestId(requestId); packet.setBizData(request); exchange.getIn().setBody(packet.serialize()); }
第三步:配置聚合器实现请求响应匹配
使用Camel聚合EIP暂存请求上下文,收到响应后通过ID匹配,自动返回给对应的REST调用方:
// 1. 声明聚合仓库,用于暂存请求上下文(生产环境可替换为Redis等持久化仓库) @Bean public AggregationRepository tcpAggregationRepository() { return new MemoryAggregationRepository(); } // 2. 原有REST接收路由改造 rest() .consumes("application/json").produces("application/json") .post("/tcp") .type(RequestBean.class) .route() .process(this::transformTcpMessage) // 按requestId聚合,暂存请求上下文 .aggregate(header("TCP_REQUEST_ID"), (oldExchange, newExchange) -> { if (oldExchange == null) { // 首次请求进来,暂存原请求上下文 return newExchange; } // 响应回来,将响应数据写入原请求上下文 oldExchange.getIn().setBody(newExchange.getIn().getBody()); return oldExchange; }) .aggregationRepository(tcpAggregationRepository()) // 配置30秒超时,避免无效请求长期占用资源 .completionTimeout(30000) // 异步发送TCP请求,发送后进入聚合等待状态 .to("netty:tcp://127.0.0.1:9898?sync=false&keepAlive=true&tcpNoDelay=true&reuseAddress=true&clientMode=true") // 聚合完成后自动将响应返回给REST调用方 .endRest(); // 3. 单独配置TCP响应监听路由 from("netty:tcp://127.0.0.1:9898?sync=false&keepAlive=true&clientMode=true") // 解析响应报文,提取关联的requestId存入Header .process(exchange -> { byte[] respBytes = exchange.getIn().getBody(byte[].class); TcpRespPacket resp = TcpRespPacket.deserialize(respBytes); exchange.getIn().setHeader("TCP_REQUEST_ID", resp.getRequestId()); exchange.getIn().setBody(resp.getBizData()); }) // 将响应送入聚合器匹配对应请求 .aggregate(header("TCP_REQUEST_ID"), new UseLatestAggregationStrategy()) .aggregationRepository(tcpAggregationRepository()) .completionSize(1) .end();
注意事项
- 必须和TCP服务端约定请求响应统一携带
requestId,这是实现匹配的核心基础 - 高并发场景建议使用持久化聚合仓库,避免服务重启导致请求上下文丢失
- 单长连接并发不足时,可以配置Netty连接池参数开启多长连接复用,不影响请求匹配逻辑
内容的提问来源于stack exchange,提问作者Pubudu Ranathunga
相关产品推荐
相关产品推荐

