当前位置: 首页 > news >正文

手写 RPC 框架零拷贝实战:把 Codec 从 byte[] 搬到 ByteBuf,一次干掉全链路内存拷贝

本文以开源 RPC 框架 jaws(Java 17 + Netty,目标是用 1.7 万行代码吸纳 Dubbo 的核心能力)最近一轮编解码层重构为背景,自底向上拆解一次 RPC 请求的编码、传输、解码全过程中每一次内存分配和拷贝发生在哪里、为什么发生、怎么消除,以及零拷贝改造引入的新问题——引用计数泄漏的排查与修复。

〇、背景:一次对比分析暴露的效率短板

jaws 的定位是"保持轻量的同时吸纳 Dubbo 的核心能力"。在完成传输层瘦身、Directory/Filter 对齐 Dubbo 等三轮重构后,我写了一份 codec 层的 Dubbo 对比分析文档,把 jaws 和 Dubbo 的编解码路径逐次分配/拷贝数了一遍,结论不太好看:

Jaws 旧编码路径: JawsCodec.encodeRequest() ① ByteArrayOutputStream → body byte[] [分配 + 拷贝] ② encode() → new byte[16] header + arraycopy 拼接 [分配 + 拷贝] ③ encode() → new byte[16 + bodyLen] + 两次 arraycopy [分配 + 2 拷贝] FrameEncoder.encodeFrame() ④ new byte[16 + total] → arraycopy 嵌入传输帧 [分配 + 拷贝] NettyEncoder.encode() ⑤ ByteBuf.writeBytes(finalFrame) [拷贝到 ByteBuf]

不算序列化器内部的缓冲,一次请求编码至少 5 次分配、5 次拷贝。而 Dubbo 的ExchangeCodec是在 ChannelBuffer 上预留 header 空间、body 直写、最后回填 header,中间 byte[] 分配接近于零。

差距的根源不在写法,而在接口模型:jaws 的 Codec SPI 签名是byte[] encode(Channel, Object)——接口把"byte[]"定为契约,那么每个跨层边界都必然产生一次数组拷贝。所以这轮重构的思路很明确:改契约,而不是改实现。顺带把对比文档里列的另一个改进项——serializationId 嵌入协议头——也一起落地了。

一、先搞清楚:旧架构的拷贝为什么会发生

自底向上看旧编解码的分层:

NettyEncoder/Decoder (传输层 - 传输帧拆包/组帧, magic=0xF1F1) ↕ FrameEncoder (帧封装层 - 把协议帧嵌入传输帧) ↕ JawsCodec (协议层 - 业务编解码, magic=0xF0F0)

这是一个双 magic 双层帧结构:JawsCodec 产出 16B 协议头 + body 的 byte[],FrameEncoder 在外面再包一层 16B 传输头(0xF1F1),最后 NettyEncoder 把整个 byte[] 写进 Netty 的 ByteBuf。

每一层拿到的都是上一个层的byte[],而每层的产出又必须是新的byte[](因为要加自己的头),于是 arraycopy 不可避免。更隐蔽的浪费有两个:

一是ByteArrayOutputStream.toByteArray()本身就是一次内部缓冲数组的防御性拷贝;二是序列化每个参数时serialize()返回独立的 byte[] 再被 ObjectOutputStream 包装一层,多参数请求的分配次数随参数个数线性增长。

解码侧同理:NettyDecoderin.readBytes(data)把完整协议帧拷进 byte[],JawsCodec 再System.arraycopy提取 body,然后包 ByteArrayInputStream 反序列化——又是两次分配两次拷贝。

结论:只要 Codec 的契约是 byte[],拷贝次数就是分层层数 × 2 起步。要根治,就得让 Codec 直接操作 Netty 的 ByteBuf,让"最终目的缓冲区"贯穿所有层。

二、改造一:编码路径——预留 header、直写 body、回填

新接口:

@Spi(scope=Scope.PROTOTYPE)publicinterfaceCodec{voidencode(Channelchannel,Objectmessage,ByteBufout)throwsIOException;Objectdecode(Channelchannel,ByteBufin)throwsIOException;}

调用方NettyChannel直接把 Netty 分配的目标缓冲区传进来:

ByteBufbuf=channel.alloc().buffer();try{codec.encode(this,request,buf);}catch(IOExceptione){buf.release();thrownewJawsServiceException("encode request error: url="+getUrl().getUri(),e);}ChannelFuturewriteFuture=channel.writeAndFlush(buf);

注意两个细节:channel.alloc().buffer()用的是该 channel 绑定的分配器(走 Netty 的池化 PooledByteBuffer,本身就是省 GC 的关键);编码失败时主动release(),因为writeAndFlush不会发生,Netty 不会替你释放。

encodeRequest的核心写法——这正是 DubboExchangeCodec的经典手法:

privatevoidencodeRequest(Channelchannel,Requestrequest,ByteBufout)throwsIOException{Serializationserialization=ExtensionLoader.getExtensionLoader(Serialization.class).getExtension(channel.getUrl().getParameter(URLParamType.serialization));// Reserve header spaceintheaderStart=out.writerIndex();out.writerIndex(headerStart+HEADER_LENGTH);// Write body directly to ByteBufByteBufOutputStreambodyOut=newByteBufOutputStream(out);ObjectOutputoutput=createOutput(bodyOut);output.writeUTF(request.getInterfaceName());output.writeUTF(request.getMethodName());output.writeUTF(request.getParamDesc());if(request.getArguments()!=null){for(Objectobj:request.getArguments()){serialize(output,obj,serialization);}}// ... attachments ...output.flush();output.close();intbodyLength=out.writerIndex()-headerStart-HEADER_LENGTH;// Backfill header with serializationId embedded in flagbyteflag=(byte)(FLAG_REQUEST|((serialization.getSerializationNumber()<<3)&SERIALIZATION_MASK));writeHeader(out,headerStart,flag,request.getRequestId(),bodyLength);}

三步走:先把 writerIndex 前推 16 字节留出 header 空间;body 通过ByteBufOutputStream直写进最终缓冲区(ObjectOutputStream 直接建立在它之上,中间零中转);写完后用out.writerIndex() - headerStart - HEADER_LENGTH算出 body 长度,回填 header。

body 从序列化器到网络发送只经过一个池化 ByteBuf,中间 byte[] 分配清零。body 长度只有写完才知道,所以"预留 + 回填"是唯一顺序,这也是为什么协议头里必须有 bodyLen 字段——协议设计决定了编码实现能不能做零拷贝。

服务端的sendResponse同构,同样直写 channel 的 alloc 分配的缓冲区。

三、改造二:解码路径——readRetainedSlice 切片零拷贝

传输层的NettyDecoder(继承ByteToMessageDecoder)完成拆包后,不再把帧拷进 byte[],而是直接把 ByteBuf 切片传下去:

// Pass ByteBuf directly to JawsCodec.decode (zero-copy, no frame byte[] allocation)in.resetReaderIndex();// Retain the buffer since the caller (ByteToMessageDecoder pipeline) may release it;// NettyChannelHandler is responsible for releasing after processing.ByteBufframe=in.readRetainedSlice(JawsCodec.HEADER_LENGTH+bodyLength);NettyMessagemessage=newNettyMessage(isRequest,requestId,frame);out.add(message);

readRetainedSlicereadSlice的区别是前者会把引用计数 +1——这个细节是后文泄漏问题的伏笔。NettyMessage也随之改成了携带 ByteBuf 的 record:

publicrecordNettyMessage(booleanisRequest,longrequestId,ByteBufdata){}

JawsCodec 解码时同样切片,body 部分用ByteBufInputStream直接包给反序列化器:

// Slice body region from ByteBuf (zero-copy, no body byte[] allocation)ByteBufbodyBuf=in.retainedSlice(in.readerIndex(),bodyLength);in.skipBytes(bodyLength);try(ByteBufInputStreambodyIn=newByteBufInputStream(bodyBuf)){ObjectInputinput=createInput(bodyIn);if(isResponse){returndecodeResponse(input,dataType,requestId,serializationId,serialization);}else{returndecodeRequest(input,requestId,serializationId,serialization);}}finally{bodyBuf.release();}

解码路径上 body 也没有任何中间数组——反序列化器直接从网络缓冲区的切片视图里读。旧路径的"readBytes 拷一次 + arraycopy 拷一次"归零。

四、改造三:顺手砍掉双 magic 传输帧

接口 ByteBuf 化之后回头看旧分层,FrameEncoder 那层的存在理由已经不成立了。它当年干的活是:校验 JawsCodec 输出的 magic、套 16B 传输头(0xF1F1)、做编码降级。而现在帧检测由 NettyDecoder 直接校验协议 magic(0xF0F0)完成,编码降级挪进了 sendResponse 的异常处理。

于是这轮把FrameEncoder(80 行)和 NettyEncoder(32 行)整个删掉,协议变成单层 16B header:

Bytes 0-1 : magic 0xF0F0 Byte 2 : version (当前 = 1) Byte 3 : flag (低 3 位 = 消息类型, 高 5 位 = serializationId) Bytes 4-11 : requestId Bytes 12-15 : body length

收益是三重的:每个请求省 16 字节传输帧头、省一次帧封装的分配拷贝、Netty pipeline 少一跳 handler(MessageToByteEncoder内部还有一次请求跨线程的 Promise 语义开销)。

有人会问:双层帧不是为"多协议复用传输层"预留的吗?实践中的判断是:YAGNI。传输帧的唯一消费者就是 jaws 协议自己,这层抽象没有第二个使用者时,它就是纯税。Dubbo 的单层 0xdabb 帧跑了十年也证明 16B header 够用。

五、改造四:serializationId 嵌入 flag 高 5 位——每消息独立序列化

这是 Dubbo 协议早就有、jaws 一直缺的能力:序列化方式跟着消息走,而不是跟着连接配置走

Dubbo 把 serializationId 放在 flag 的低 5 位。jaws 的 flag 低 3 位已被消息类型占用(0x00=request / 0x01=response / 0x03=void / 0x05=exception),高 5 位正好空闲,语义完全同构:

实现上分三步。第一步给 SPI 注解加数字标识:

public@interfaceSpiMeta{Stringname();/** Optional numeric identifier for the SPI extension... Value must be in range 0-31. Default -1 means not assigned. */intnumber()default-1;}
@SpiMeta(name="hessian2",number=0)publicclassHessian2SerializationextendsAbstractSerialization{...}@SpiMeta(name="fastjson2",number=1)publicclassFastJson2SerializationextendsAbstractSerialization{...}

第二步 ExtensionLoader 建一张 number → name 的映射,O(1) 反查:

publicTgetExtensionByNumber(intnumber){checkInit();Stringname=numberToName.get(number);returnname!=null?getExtension(name):null;}

第三步是全链路闭环,这是最容易漏的部分。编码端在第二节已经看到:serialization.getSerializationNumber() << 3写进 flag。解码端提取:

byteflag=in.readByte();bytedataType=(byte)(flag&MASK);byteserializationId=(byte)((flag&SERIALIZATION_MASK)>>3);

然后request 上带着它进入业务处理,handler 把它拷到 response,encodeResponse 据此选序列化器

// NormalRequestHandlerresponse.setSerializationNumber(request.getSerializationNumber());
privatevoidencodeResponse(Responseresponse,ByteBufout)throwsIOException{// Use the serialization carried on the response (copied from request by handler),// so the response is encoded with the same serialization the client used.Serializationserialization=ExtensionLoader.getExtensionLoader(Serialization.class).getExtensionByNumber(response.getSerializationNumber());...}

为什么 response 必须跟随 request?因为消费端解码用的是自己配置的序列化方式,如果服务端用配置里的另一种方式编码响应,消费端就解不出来了。这个"请求-响应序列化一致性"约束,就是 Dubbo 在协议头里放 serializationId 的根本原因。落地之后,同一服务端可以同时接收 hessian2 客户端和 fastjson2 客户端的调用,互不干扰——多客户端异构场景打通。

六、零拷贝的代价:两处引用计数泄漏

Byte[] 是 JVM 管理的,扔了就有 GC 兜底;ByteBuf 是池化 + 引用计数的,谁 retain 谁 release,少一次 release 就泄漏一块池化内存,而且 GC 救不了。这正是零拷贝改造最需要警惕的地方——我这轮 review 就抓出了两处。

Netty 的引用计数规则先说清楚:readRetainedSlice/retainedSlice会 +1;ByteToMessageDecoder在一次channelRead结束后释放自己持有的引用;所以只要切片对象要在 channelRead 返回后继续存活(比如提交给业务线程池),就必须再 retain,并在用完后 release。

泄漏点一:线程池拒绝路径。NettyChannelHandler.channelRead里,异步路径先 retain 再提交线程池:

if(threadPoolExecutor!=null){try{// Retain the ByteBuf for async processing (pipeline may release after this method returns)nettyMsg.data().retain();threadPoolExecutor.execute(()->{try{processMessage(ctx,nettyMsg);}finally{nettyMsg.data().release();}});}catch(RejectedExecutionExceptionrejectException){// Only server-side requests go through the thread pool;// reject and return error response to client when pool is full.rejectMessage(ctx,nettyMsg);// ← 旧代码:这里没 release!}}

任务正常执行时:异步任务的 finally release 一次 + 外层 finally release 一次 = 两次,抵消 retain 和 readRetainedSlice 各一次,平衡。但execute()RejectedExecutionException时任务根本没提交成功,异步 release 不会发生——外层 finally 只 release 掉 readRetainedSlice 那次,retain 的那次就漏了。最坑的是触发条件是线程池打满,恰好是流量高峰,泄漏速度和事故规模正相关。修复就是 catch 里补一次:

}catch(RejectedExecutionExceptionrejectException){// Release the ByteBuf retained above; the finally block below only releases// the reference acquired by the decoder (readRetainedSlice).nettyMsg.data().release();rejectMessage(ctx,nettyMsg);}

泄漏点二:异常抛出早于 try-with-resources。JawsCodec.decode 里,旧顺序是先切片、后校验 serializationId:

// 旧顺序(有泄漏风险)ByteBufbodyBuf=in.retainedSlice(in.readerIndex(),bodyLength);// +1in.skipBytes(bodyLength);Serializationserialization=...getExtensionByNumber(serializationId);if(serialization==null){thrownewJawsFrameworkException("decode error: unknown serializationId "+serializationId,...);// ← bodyBuf 已经 +1,但 try-with-resources 还没进来,永远没人 release}try(ByteBufInputStreambodyIn=newByteBufInputStream(bodyBuf)){...}

retainedSlice和 try 块之间任何一行抛异常,这次引用计数就丢了。恶意客户端持续发 unknown serializationId 的请求,可以稳定放大泄漏。修复方式是把纯校验逻辑挪到任何切片之前——“先做所有不持有资源的校验,再获取资源”,这是引用计数环境下的通用防御性顺序

// Resolve serialization from the id embedded in the protocol header// (must happen before the retainedSlice below, otherwise an unknown id// would throw and leak the sliced body buffer)Serializationserialization=ExtensionLoader.getExtensionLoader(Serialization.class).getExtensionByNumber(serializationId);if(serialization==null){thrownewJawsFrameworkException(...);}ByteBufbodyBuf=in.retainedSlice(in.readerIndex(),bodyLength);in.skipBytes(bodyLength);try(ByteBufInputStreambodyIn=newByteBufInputStream(bodyBuf)){...}finally{bodyBuf.release();}

顺带一提:ByteBufInputStream.close()不会释放包装的 ByteBuf,所以 finally 里的显式 release 不能省——这是 Netty 的一个经典反直觉点,很多人以为 try-with-resources 能兜底,其实不能。

排查这类问题的标准姿势是开 Netty 的泄漏检测压测:

# SIMPLE 是采样检测,PARANOID 是全量检测(性能损失大,仅排查用)java-Dio.netty.leakDetection.level=PARANOID-jarprovider.jar

泄漏发生时会打出LEAK: ByteBuf.release() was not called before it's garbage-collected的堆栈,指明 retain 的来源。

七、收益结算与 Dubbo 对照

这轮重构整体 56 文件 +1003/-498(同期还有 jaws-transport-netty 模块并入 jaws-core 等结构调整),编解码路径的前后对照:

维度改造前改造后
编码中间分配≥5 次(body/header/拼接/传输帧/writeBytes)0 次(池化 ByteBuf 直写)
编码拷贝≥5 次0 次
解码中间分配2 次(帧 byte[] + body byte[])0 次(retainedSlice 视图)
每请求协议头开销32B(传输帧 16B + 协议帧 16B)16B
序列化方式选择按 URL 配置,连接级flag 高 5 位,消息级,响应自动跟随请求
内存管理责任GC 兜底引用计数,需显式 release(新增两类泄漏风险点)

对照 Dubbo:协议层的序列化透传、零拷贝编码、预留回填手法已经完全同构;jaws 仍缺的是 FLAG_EVENT 心跳位、twoway/oneway 标记(flag 还有 bit7/bit6 两个空位,实现方式和 serializationId 完全同构,属于低成本跟进项)和 status 状态码粒度。

最后说点务虚的。这轮改造给我的最大触动是:性能问题的根因经常在接口契约层,而不是实现层。当 Codec 的签名是byte[]时,写得多聪明都躲不开边界拷贝;把签名改成ByteBuf,正确实现自然就是零拷贝。但换契约是有代价的——byte[] 世界里 GC 兜底的一切,到引用计数世界里都要人工对账,retain/release 的配对正确性成了新的正确性维度。零拷贝从来不是免费的,它只是把成本从运行时延迟转移到了工程纪律上。


jaws 项目地址:github.com/javahongxi/jaws,一个用 1.7 万行代码吸纳 Dubbo 核心能力的轻量 RPC 框架,支持 ZK/Nacos 注册中心、五种负载均衡、三种容错策略、泛化调用、动态配置、MCP 桥接。编解码设计详见 doc/codec-comparison.md。

http://www.cnnetsun.cn/news/4152450.html

相关文章:

  • 浏览器下载速度慢的成因分析与全链路优化指南
  • 数学建模中变量区分度分析:t检验、点二列相关与Cronbach‘s Alpha实战指南
  • AI绘图实战:用提示词工程为电商产品批量生成高转化率视觉素材
  • C# TCP/IP网络编程实战:从Socket到健壮通信框架
  • HLSL程序化砖墙材质:从数学逻辑到虚幻引擎实战
  • 数模实战中的描述分析内功:从数据诊断到建模决策
  • 微信小程序用户信息获取:从wx.getUserInfo到wx.getUserProfile的完整实践指南
  • SpringBoot+Vue招聘系统开发实战与架构解析
  • 数据可视化实战:折柱混合图在数据聚合与对比分析中的应用
  • 3970亿参数大模型量化实战:NVIDIA Model Optimizer核心原理与避坑指南
  • npm与npx深度解析:从包管理到命令执行的Node.js生态核心工具
  • 大模型智能体高效压缩:知识蒸馏与模型剪枝量化实战指南
  • YOLOv5网络架构详解与实战:从原理到RV1106/RK3568部署
  • YOLOv8低光照目标检测失效机理与全链路优化
  • Vivado 2018.3安装与启动疑难排查:从环境配置到深度修复
  • 系统动力学与智能体建模:数学建模如何破解非法野生动物贸易难题
  • HTTP 5xx服务器错误全解析:从原理到实战排查与防御
  • TSAssistant:基于智能体框架与人在回路的自动化安全评估系统设计与实践
  • 数学建模Prompt设计:原子化拆解与可验证指令链
  • 大模型面试核心挑战与工程实践指南
  • 数学建模四要素:思路·模型·代码·论文的工程化方法论
  • Codex Memory Trim:解决AI命令行工具内存膨胀的维护利器
  • InternLM+Lagent+Streamlit大模型交互骨架实战
  • 层次分析法实战:从主观判断到科学权重的多准则决策指南
  • 层次分析法(AHP)实战指南:从数学建模到科学决策
  • Go应用零码改造接入观测云:OpenTelemetry与DataKit实战指南
  • 自适应滤波:时间序列动态建模的实时校准核心方法
  • SpringBoot+Vue构建高校实习就业全流程管理系统
  • 计算机网络面试核心知识:TCP/IP协议与HTTP/HTTPS详解
  • 5分钟理顺Windows右键菜单:ContextMenuManager完整指南