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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 07:57:03