尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Dubbo自定义协议实战:从SPI机制到编解码实现全解析

Dubbo自定义协议实战:从SPI机制到编解码实现全解析 自定义协议这条路Dubbo 比你想的好走做一个 RPC 框架的二次开发最容易被问倒的问题就是能不能不按默认协议走我遇到过不少团队系统和 Dubbo 深度绑定但是底层通信格式必须换成自己内部统一的一套私有协议或者说需要直接兼容老系统的报文结构。这时候如果对 Dubbo 的扩展机制没有概念很容易陷入“改源码”的恐惧里。实际上Dubbo 从设计之初就预留了协议扩展位基于 Dubbo SPI 机制实现自定义通信协议并没有想象中那么难但也绝不是改个配置就能蒙混过关的事。今天就把我实际做过的自定义协议开发过程拆开揉碎讲一遍包括扩展机制的原理、编码解码的实现、服务端和消费端接入的细节以及几个藏得挺深的坑。1. 为什么需要自研协议默认协议不是够用吗1.1 Dubbo 默认协议的核心优势与边界Dubbo 自带的dubbo://协议基于 TCP 传输配合 Hessian2 序列化内部封装了连接管理、心跳检测、粘包拆包、请求响应关联这些复杂逻辑。大多数业务场景下拿过来直接用完全没问题。默认协议厉害的地方在于它把 RPC 调用过程中的所有横切面都考虑进去了从连接池到线程派发从超时控制到泛化调用开发者不需要关心底层任何细节。但它的边界也很明显。默认协议的消息头设计是固定的请求体、响应体的结构绑定 Dubbo 自己的模型。如果你面对的是下面几类场景默认协议就会变成一种束缚公司内部存在一套已经跑了很多年的老系统通信格式是老旧的文本协议或自定义二进制格式现在要打通 Dubbo 服务与新系统之间的调用不能要求老系统改协议。开发语言异构服务端是 Java但客户端可能是 C、Python、Go直接跑 Dubbo 协议需要维护多语言 SDK成本高。需要对报文的传输体积做极致优化默认协议头里塞了太多你不需要的字段。网关层需要根据报文内容做路由或审计必须能读懂原始数据包。我从实际经验出发告诉你当出现这些情况时与其在业务代码里拦截流量做协议转换不如直接在 Dubbo 协议层把自定义协议实现出来。这样对上层业务完全透明调用方只依赖接口不感知通信细节是一种更干净的架构。1.2 自定义协议解决的真实问题有一次我遇到的项目服务端有一套内部网关所有后端服务之间的通信统一走一个类 JSON 的文本协议每条消息以换行符分隔按 key-value 组织字段。新上的 Dubbo 服务要想被老系统直接调用最简单的做法就是实现一个 Dubbo 自定义协议让 Dubbo 服务端能够接受并解析那种类 JSON 格式的报文同时把响应也编码成同样格式。核心工作就是在编解码器里做报文格式转换业务代码一行不用动。另一个场景是内部平台要对所有 RPC 流量做定制化的链路追踪默认协议解析起来比较费劲自定义协议可以在报文头里预留 traceId 字段直接在解码层就把链路信息提取出来性能比从附件里找快得多。所以自定义协议的本质不是炫技而是把传输格式的控制权拿回自己手里。2. 动手前必须先吃透 Dubbo SPI 扩展机制2.1 Dubbo SPI 与 Java SPI 的关键差异如果直接去改 Dubbo 源码那基本上是把升级路线堵死了。正确姿势是利用 Dubbo 的扩展机制在不动框架代码的前提下注入自己的协议实现。Dubbo 的扩展机制基于一套增强版的 SPIService Provider Interface。和 Java 自带的 SPI 相比它有几个核心改进每个扩展点可以有默认实现通过注解SPI(dubbo)声明默认扩展名。同一个扩展点可以存在多个实现按名称区分通过 URL 中的参数动态选择。扩展实现支持自动包装Wrapper类似 AOP 的增强。支持自动装配扩展实现里可以注入其他扩展点。支持自适应扩展通过Adaptive注解在运行时根据参数选择具体实现。Dubbo 的扩展点配置文件放在META-INF/dubbo/目录下还有META-INF/dubbo/internal/和META-INF/services/文件名为扩展点接口的全限定名文件内容每一行是一个扩展名和实现类的映射。想要开发自定义协议最核心的扩展点是org.apache.dubbo.rpc.Protocol。这个接口定义了服务暴露和服务引用的两个核心方法。SPI(dubbo) public interface Protocol { int getDefaultPort(); Adaptive T ExporterT export(InvokerT invoker) throws RpcException; Adaptive T InvokerT refer(ClassT type, URL url) throws RpcException; void destroy(); }Adaptive注解在这里非常关键。它表示export和refer方法在运行时生成的适配器类会根据 URL 中的protocol参数找到对应的扩展实现。也就是说当注册中心 URL 的协议头是x://时Dubbo 会自动路由到名为x的 Protocol 实现。2.2 搞懂 Protocol 和 Remoting 层的关系Protocol 是 Dubbo RPC 分层中的最上层抽象它本身不直接处理 TCP 连接和报文编解码而是委托给更底层的 Remoting 模块。Remoting 层负责传输、编码解码、心跳、连接管理等核心扩展点包括Codec2负责消息的编码和解码。Transporter负责服务端的启动和客户端的连接默认基于 Netty 实现。Dispatcher负责消息派发策略。ThreadPool负责业务线程池。自定义协议时比较常规的方案是Protocol 实现类负责任务调度和服务管理底层的 TCP 传输和编解码复用 Dubbo 的 Remoting 层只需要实现自己的 Codec2。当然如果你的协议底层不走 TCP而是走 HTTP 或者 UDP那么 Transporter 也需要一起扩展但这是少数情况。我把协议扩展的几个层次列一下你就理解工作量了扩展目标扩展接口改动范围只改报文格式Codec2中改报文格式 请求响应模型Codec2 Protocol中改传输层Codec2 Protocol Transporter大改序列化方式Serialization小大部分团队自定义协议的实际需求停留在改报文格式所以重点花在 Codec2 和 Protocol 上就对了。3. 从零写一个自定义协议X 协议的完整实现3.1 协议头设计动手写代码之前先把协议格式定下来。我记得当初设计的时候参考了 Dubbo 默认协议头又结合了自己的业务需求。定义一个相对简单但完整的二进制协议命名为x协议头固定 16 字节魔数4 字节用于快速校验协议合法性定义为0x58425250对应 ASCII 码 XBRP。协议版本1 字节方便后续协议升级兼容。消息类型1 字节区分请求、响应、心跳也方便后续扩展。请求 ID8 字节long 型用于关联请求和响应也用来做日志追踪。数据长度2 字节无符号 short表示 body 的字节长度。为什么用 16 字节的头部因为 TCP 是流式传输接收方必须通过固定长度的头来确认读多少字节算一条完整消息这是所有二进制自定义协议的基础。头部固定长度之后解码器就能够做到先收头、再收体。为什么不把魔数设计成 2 字节魔数的作用是快速识别协议降低误判概率。4 字节的魔数在二进制流中出现巧合的概率远低于 2 字节尤其当服务端可能同时存在多个协议端口时4 字节能显著减少解析错乱的可能。3.2 扩展点声明与 Protocol 实现类先加入 Dubbo 依赖。这里用 Apache Dubbo 3.x 版本作为示例dependency groupIdorg.apache.dubbo/groupId artifactIddubbo/artifactId version3.2.0/version /dependency dependency groupIdorg.apache.dubbo/groupId artifactIddubbo-remoting-netty4/artifactId version3.2.0/version /dependency在META-INF/dubbo/org.apache.dubbo.rpc.Protocol文件中写入xcom.example.dubbo.protocol.XProtocol然后实现 Protocol 接口。不建议直接实现裸接口可以继承AbstractProtocol它已经把exporterMap、serverMap这些通用状态管理好了public class XProtocol extends AbstractProtocol { public static final String NAME x; Override public T ExporterT export(InvokerT invoker) throws RpcException { URL url invoker.getUrl(); String key ServiceKey.key(url.getPort(), url.getPath(), url.getVersion()); XExporterT exporter new XExporter(invoker, key, exporterMap); exporterMap.put(key, exporter); openServer(url); return exporter; } Override public T InvokerT refer(ClassT type, URL url) throws RpcException { XInvokerT invoker new XInvokerT(type, url, url.getPath()); invokers.add(invoker); return invoker; } private void openServer(URL url) { String key url.getAddress(); boolean isServer url.getParameter(server, true); if (isServer) { RemotingServer server serverMap.get(key); if (server null) { server new XServer(url, new XCodec2()); serverMap.put(key, server); } } } }这里有个设计细节很关键多个服务暴露在同一个 IP:Port 时serverMap 会复用同一个 RemotingServer。如果不做这个复用每暴露一个服务就启动一个新端口资源很快就会被打满。3.3 编解码器的实现核心中的核心编解码器是自定义协议最重要的部分。Dubbo 的Codec2接口包含encode和decode两个方向encode 是把响应或请求对象变成字节流decode 是把字节流变回对象。public class XCodec2 implements Codec2 { private static final int HEADER_LENGTH 16; private static final int MAGIC 0x58425250; Override public void encode(Channel channel, ChannelBuffer buffer, Object message) throws IOException { byte[] body serialize(message); buffer.writeInt(MAGIC); buffer.writeByte(1); // version buffer.writeByte(messageType(message)); buffer.writeLong(0); // requestId实际使用时会从 message 中提取 buffer.writeShort(body.length); buffer.writeBytes(body); } Override public Object decode(Channel channel, ChannelBuffer buffer) throws IOException { int readable buffer.readableBytes(); if (readable HEADER_LENGTH) { return DecodeResult.NEED_MORE_INPUT; } buffer.markReaderIndex(); int magic buffer.readInt(); if (magic ! MAGIC) { throw new IOException(Invalid magic number: Integer.toHexString(magic)); } buffer.readByte(); // version byte type buffer.readByte(); long requestId buffer.readLong(); int length buffer.readUnsignedShort(); if (length 0) { return null; } if (buffer.readableBytes() length) { buffer.resetReaderIndex(); return DecodeResult.NEED_MORE_INPUT; } byte[] body new byte[length]; buffer.readBytes(body); return deserialize(type, requestId, body); } }核心逻辑就是前面提到的先收头再收体判断可读字节是否够一个协议头不够就返回NEED_MORE_INPUT。读魔数、校验版本。根据头部声明的 body 长度判断是否收到完整消息。数据不完整时重置读索引等待下次读取。数据完整时读取 body 并反序列化为业务对象。解码器的返回对象可以是DecodeResult.NEED_MORE_INPUT、null心跳或业务消息对象。Dubbo 的 Remoting 层会循环调用 decode直到不能再解析出完整消息。这里我想强调一个很多人会犯的错忘记调用buffer.markReaderIndex()或者不好好利用它的 reset 能力。网络包是流式的一次decode调用拿到的数据可能只是半个包也可能是好几个包粘在一起。只有在数据不完整时把读索引重置到包头位置才能保证下个包到达后从正确的位置开始解析。3.4 XInvoker 与请求派发Protocol.refer返回的 Invoker 是消费端发起调用的入口。XInvoker需要实现doInvoke方法把 Dubbo 的RpcInvocation转成自定义协议的消息对象通过网络发送给服务端并同步等待响应public class XInvokerT extends AbstractInvokerT { private final String path; private final RemotingClient client; private final ConcurrentMapLong, CompletableFutureObject pending new ConcurrentHashMap(); public XInvoker(ClassT type, URL url, String path) { super(type, url, new String[]{INTERFACE_KEY, GROUP_KEY, TOKEN_KEY, TIMEOUT_KEY}); this.path path; this.client new XClient(url, new XCodec2()); } Override protected Result doInvoke(Invocation invocation) throws Throwable { RpcInvocation inv (RpcInvocation) invocation; long requestId REQUEST_ID_GENERATOR.incrementAndGet(); XRequest request new XRequest(requestId, path, inv.getMethodName(), inv.getArguments()); CompletableFutureObject future new CompletableFuture(); pending.put(requestId, future); client.send(request); return (Result) future.get(); } }服务端接收消息后需要从请求中解析出接口名、方法名、参数再通过本地暴露的 Invoker 去反射调用。这是自定义协议开发中比较繁琐的部分但好在 Dubbo 提供了现成的DubboProtocol可以参考它内部通过requestHandler把远程请求和本地 exporter 关联起来。服务端处理的核心逻辑public void receive(Channel channel, Object message) { if (message instanceof XRequest) { XRequest request (XRequest) message; String serviceKey ServiceKey.key(channel.getUrl().getPort(), request.getPath(), null); Exporter? exporter exporterMap.get(serviceKey); Invoker? invoker exporter.getInvoker(); RpcInvocation invocation new RpcInvocation( request.getMethodName(), request.getParameterTypes(), request.getArguments()); Result result invoker.invoke(invocation); channel.send(new XResponse(request.getRequestId(), result.getValue(), result.getException())); } else if (message instanceof XResponse) { XResponse response (XResponse) message; CompletableFutureObject future pending.remove(response.getRequestId()); if (future ! null) { future.complete(response.getResult()); } } }3.5 请求还是响应消息类型别搞混自定义协议的消息类型字段type设计为 4 种枚举值就够了请求1、响应2、心跳请求3、心跳响应4。心跳消息不需要 body或只需要很小的 body。实现心跳能力对于保持长连接非常重要否则空闲连接会在网络设备上被回收。Dubbo 的HeaderExchangeHandler里本来就内置了心跳机制但因为是自定义协议需要自己在编解码层处理心跳消息或者在Codec2里拦截处理。我在做的时候心跳逻辑直接放在服务端的 receive 方法里收到心跳请求就回一个心跳响应收到心跳响应就把对应的时间戳更新一下用于后续的可用性判断。4. Provider 和 Consumer 接入怎么把自定义协议跑起来4.1 服务端暴露服务写完扩展之后接入方式其实非常简单。Provider 端的 XML 配置里把协议名从dubbo换成x。dubbo:application namex-provider / dubbo:registry addressnacos://127.0.0.1:8848 / dubbo:protocol namex port20880 / dubbo:service interfacecom.example.DemoService refdemoService /或者用注解方式在服务实现类上直接指定协议Service(protocol x) public class DemoServiceImpl implements DemoService { // ... }启动服务后注册中心上出现的 URL 会是x://192.168.1.100:20880/com.example.DemoService?...。Dubbo 的 Registry 模块不会关心协议名是什么它只负责把这个 URL 保存到注册中心消费端拉取到 URL 后再根据protocolx找到对应的 Protocol 实现。这里特别强调一点协议的绑定发生在 URL 层面而不是代码层面。所以只要 Protocol 扩展实现正确Provider 端就完全没有其他额外代码。4.2 消费端引入服务Consumer 端配置同样简单dubbo:application namex-consumer / dubbo:registry addressnacos://127.0.0.1:8848 / dubbo:reference iddemoService interfacecom.example.DemoService protocolx /或者注解方式DubboReference(protocol x) private DemoService demoService;消费端启动时会从 Nacos 拉取x://协议的 URLDubbo 的ReferenceConfig会触发Protocol.refer的调用从而创建自定义协议的 Invoker。整个过程对业务代码完全透明调用方只需要依赖接口不需要关心底层走的是什么协议。4.3 注册中心选型Nacos 与自定义协议的配合自定义协议和注册中心之间解耦得相当彻底。只要 URL 中的协议字段能被注册中心正确存储那么 Nacos、Zookeeper 都没有区别。我在项目里用的是 Nacos原因无非是公司内部基础设施统一了。需要注意的只有一点Nacos 注册元数据时默认有preserved.heart.beat.timeout这些参数如果自定义协议的心跳间隔和注册中心的心跳参数配置差距太大可能导致服务实例被误判为不健康。所以自定义协议的心跳间隔建议设置在 10 到 30 秒之间注册中心的健康检查配置需要同步调整。还有一个细节如果你使用的是 Dubbo 3.x 的应用级服务发现注册到 Nacos 的元数据会走额外的 MetadataService 接口。此时 MetaDataService 默认走dubbo协议不会自动切换成自定义协议。如果只开放了自定义协议的端口就会导致消费端拉不到元数据。解决办法是不要禁用默认协议让 MetadataService 走 dubbo 协议业务接口走自定义协议或者实现 MetadataService 时也把它注册到自定义协议上。5. 这些坑我替你先踩了一遍5.1 Wrapper 类没有加载导致增强失效Dubbo SPI 支持自动包装Wrapper比如ProtocolFilterWrapper和ProtocolListenerWrapper会在协议外层做过滤器链和监听器装配。这些 Wrapper 类的构造函数需要接收一个 Protocol 类型的参数。如果自定义 Protocol 没有正确继承AbstractProtocol或者 SPI 文件的格式写错了Wrapper 类加载失败问题不会在启动时马上暴露而是等第一个请求进来发现过滤器链完全没生效时才开始排查那时候定位成本就高了。排查方法在启动日志里搜索use dubbo protocol相关的 Wrapper 加载日志再检查Protocol$Adaptive生成的代码里是否包含 Wrapper。5.2 序列化方式不统一编解码隐藏炸雷自定义协议必然涉及序列化的选择。如果你实现了自己的 Codec2但内部使用了Hessian2Serialization而消费端用的却是FastJsonSerialization两边序列化方式不一致就会在调用时出现各种诡异异常。这里要特别提醒序列化扩展和编解码扩展是两层概念容易搞混。Codec2 负责一条消息从哪里开始到哪里结束Serialization 负责对象在 body 里如何变成字节。自定义协议可以用 Dubbo 自带的序列化扩展也可以自己在 encode/decode 方法里直接调用 Jackson、Protobuf。如果使用自带的序列化需要在 URL 上传递serialization参数URL url URL.valueOf(x://127.0.0.1:20880/com.example.DemoService?serializationfastjson2);同时你的编解码器里要根据 URL 参数动态获取序列化实现Serialization serialization ExtensionLoader.getExtensionLoader(Serialization.class) .getExtension(url.getParameter(serialization, hessian2));5.3 粘包拆包处理不当连接直接废掉这是自定义协议实现中最容易出问题的地方。TCP 是基于字节流的没有消息边界。应用层必须自己识别消息边界。我之前见过同事写的解码器没有处理一次收到多个消息的情况导致第一条消息解析成功第二条消息从错误的位置开始读Magic 直接校验失败连接被强制关闭。正确做法是使用 while 循环在decode方法里持续尝试解析try { Object msg; do { msg decodeSingle(channel, buffer); if (msg DecodeResult.NEED_MORE_INPUT) { break; } if (msg ! null) { messages.add(msg); } } while (msg ! null); } finally { // 确保循环结束时读索引指向最后一条完整消息的末尾 }这种处理方式在 Dubbo 的ExchangeCodec里有现成实现建议直接把那套逻辑吃透。5.4 线程派发策略导致业务线程阻塞自定义协议的服务端在收到请求后如果直接在 I/O 线程上执行业务逻辑遇到耗时操作就会把 Netty 的 I/O 线程占满进而拖垮整个服务端。解决方案是通过 Dubbo 的Dispatcher扩展点或ExecutorRepository来管理线程池把业务执行放到独立的业务线程池中。继承DubboProtocol里requestHandler的做法在收到请求后向线程池提交任务ExecutorService executor executorRepository.createExecutorIfAbsent(url); executor.execute(() - { Result result invoker.invoke(invocation); channel.send(new XResponse(requestId, result.getValue(), result.getException())); });线程池大小、队列类型都要根据业务量评估否则可能出现请求堆积、超时大面积失败的问题。5.5 服务分组、版本号这些参数别忘了自定义协议越是简化越容易把 Dubbo 原有的高级能力弄丢。服务分组group、版本号version这些参数如果没有在编解码时作为 key 的一部分传递就会出现服务调错的情况。如果你暴露同一个接口的多个版本那么ServiceKey.key()里必须包含version。XProtocol 的 export 方法里写死 versionnull会导致两个版本的同一个接口被同一个 key 覆盖后暴露的版本覆盖先暴露的版本。这个 bug 非常隐蔽因为大部分情况下的第一个版本是可以正常工作的。只有在发布新版本后调用方接到旧逻辑或者直接找不到服务时才会暴露。5.6 时区、字符集编码在小协议里也会咬人文本类自定义协议里时间字段的格式化和解析是最容易踩坑的。我曾经遇到过服务端用的是yyyy-MM-dd HH:mm:ss客户端传的是yyyy-MM-ddTHH:mm:ssZ两边没对齐导致所有涉及时间的字段解析全部失败。更麻烦的是这种问题在测试环境不容易暴露因为测试数据往往用的是同一个时区。建议在协议设计文档里就把时间字段的格式、时区明确写死最好统一用时间戳long。6. 从自定义协议到协议网关扩展思路再往前走一步自定义协议的实现流程跑通之后可以做很多有意思的扩展。我后来就把这套能力用在了协议网关上所有 Dubbo 服务暴露的 URL 协议头统一换成网关自定义协议网关收到报文后先解析协议头再根据接口名动态决定转发到哪个后端集群。这样网关上游的调用方和下游的 Dubbo 服务完全解耦新增服务不需要改动网关代码。对比方案双协议并存 vs 全量切换如果你担心自定义协议不够稳定可以先双协议并存。Provider 上注册两个协议dubbo 协议保留给内部系统自定义协议给新接入的调用方。注意 Dubbo 的协议并存是protocol标签里配置多个协议名或者注册多个dubbo:protocol标签。只要注册中心的 URL 是两份消费端各取所需即可。不过双协议并存会增加连接数和内存开销不建议长期使用。等自定义协议稳定运行一段时间后把 dubbo 协议的端口关掉全量切换过去。自定义协议的网关还有一个额外优势可以在协议层直接做流量复制。因为解码时能拿到完整的请求报文复制一份发给压测环境服务对业务代码完全无侵入。这在做线上全链路压测时十分有用。像这种上层架构的玩法其实都是建立在协议层有能力自由解析报文的基础上。所以说把自定义协议的实现能力掌握在手里不只是在做一个技术点而是在给整个服务治理体系打开一扇门。最后再分享一个小技巧写完自定义协议后先不要着急接入业务而是写一个简单的 echo 测试客户端和服务端各跑一个 main 方法先把一条字符串消息的编码解码调通再逐步加入 Dubbo 的接口调用逻辑。协议层的问题越早暴露越好查别等它混在业务逻辑里才来排查。
返回列表