myrpc 学习笔记-006 自定义协议
项目地址 欢迎访问
笔记总览 myrpc 学习笔记-001 实现简易版 rpc myrpc 学习笔记-002 配置加载 myrpc 学习笔记-003 Mock 服务代理 myrpc 学习笔记-004 序列化实现和 SPI 机制 myrpc 学习笔记-005 注册中心 myrpc 学习笔记-006 自定义协议 myrpc 学习笔记-007 负载均衡 myrpc 学习笔记-008 重试机制 myrpc 学习笔记-009 容错机制 myrpc 学习笔记-0010 启动机制和注解驱动
架构图v6.0.0
- 自定义协议
- 实现协议消息编解码器

一、为什么要自定义 RPC 协议?
- HTTP 协议头太重,纯 RPC 不需要那么多信息
- 自定义协议更轻量、更快、更安全
- 方便扩展序列化器、消息类型、状态码
- 解决 TCP 粘包/拆包问题
- 统一消息格式,方便客户端 + 服务端通信
二、自定义协议整体结构
协议采用 固定头 + 可变体 格式:
▼text复制代码[ 消息头(17字节 固定)] + [ 消息体(长度不固定)]
1. 消息头结构(17字节)
| 偏移 | 字段 | 长度 | 作用 |
|---|---|---|---|
| 0 | magic(魔数) | 1byte | 验证消息合法性 |
| 1 | version(版本) | 1byte | 协议升级 |
| 2 | serializer(序列化器) | 1byte | jdk/json/kryo/hessian |
| 3 | type(消息类型) | 1byte | 请求/响应/心跳 |
| 4 | status(状态) | 1byte | 成功/失败 |
| 5~12 | requestId | 8byte | 请求唯一ID |
| 13~16 | bodyLength | 4byte | 消息体长度 |
总长度:1+1+1+1+1 +8 +4 = 17byte
2. 消息体
由序列化器将 RpcRequest / RpcResponse 序列化为字节数组,长度由头中的 bodyLength 决定。
三、协议核心常量
▼java复制代码public interface ProtocolConstant { // 消息头固定长度 int MESSAGE_HEADER_LENGTH = 17; // 协议魔数 byte PROTOCOL_MAGIC = 0x1; // 协议版本 byte PROTOCOL_VERSION = 0x1; }
四、协议消息模型 ProtocolMessage
▼java复制代码@Data @AllArgsConstructor @NoArgsConstructor public class ProtocolMessage<T> { // 消息头 private Header header; // 消息体(RpcRequest / RpcResponse) private T body; @Data public static class Header { private byte magic; // 魔数 private byte version; // 版本 private byte serializer; // 序列化器 private byte type; // 消息类型 private byte status; // 状态 private long requestId; // 请求ID private int bodyLength; // 体长度 } }
五、协议枚举
1. 序列化器枚举
jdk=0,json=1,kryo=2,hessian=3
▼java复制代码public enum ProtocolMessageSerializerEnum { JDK(0, "jdk"), JSON(1, "json"), KRYO(2, "kryo"), HESSIAN(3, "hessian"); }
2. 消息状态枚举
▼java复制代码public enum ProtocolMessageStatusEnum { OK("ok", 20), BAD_REQUEST("badRequest", 40), BAD_RESPONSE("badResponse", 50); }
3. 消息类型枚举
▼java复制代码public enum ProtocolMessageTypeEnum { REQUEST(0), // 请求 RESPONSE(1), // 响应 HEART_BEAT(2), // 心跳 OTHERS(3); }
六、协议编码器 ProtocolMessageEncoder
功能
把 ProtocolMessage → 字节流 Buffer
流程
- 写入固定头(magic、version、serializer、type、status、requestId)
- 获取序列化器(根据头信息)
- 序列化消息体 →
bodyBytes - 写入bodyLength
- 写入bodyBytes
▼java复制代码public static Buffer encode(ProtocolMessage<?> protocolMessage) throws IOException { ProtocolMessage.Header header = protocolMessage.getHeader(); Buffer buffer = Buffer.buffer(); // 写入头字段 buffer.appendByte(header.getMagic()); buffer.appendByte(header.getVersion()); buffer.appendByte(header.getSerializer()); buffer.appendByte(header.getType()); buffer.appendByte(header.getStatus()); buffer.appendLong(header.getRequestId()); // 获取序列化器 Serializer serializer = SerializerFactory.getSerializer( ProtocolMessageSerializerEnum.getEnumByKey(header.getSerializer()).getValue() ); // 序列化body byte[] bodyBytes = serializer.serialize(protocolMessage.getBody()); // 写入body长度和数据 buffer.appendInt(bodyBytes.length); buffer.appendBytes(bodyBytes); return buffer; }
七、协议解码器 ProtocolMessageDecoder
功能
把 Buffer → ProtocolMessage
流程
- 读取固定头
- 校验魔数
- 获取序列化器、消息类型
- 读取bodyLength
- 截取bodyBytes
- 反序列化为
RpcRequest/RpcResponse
▼java复制代码public static ProtocolMessage<?> decode(Buffer buffer) throws IOException { // 读取头 ProtocolMessage.Header header = new ProtocolMessage.Header(); byte magic = buffer.getByte(0); if (magic != ProtocolConstant.PROTOCOL_MAGIC) { throw new RuntimeException("非法消息"); } header.setMagic(magic); header.setVersion(buffer.getByte(1)); header.setSerializer(buffer.getByte(2)); header.setType(buffer.getByte(3)); header.setStatus(buffer.getByte(4)); header.setRequestId(buffer.getLong(5)); header.setBodyLength(buffer.getInt(13)); // 解决粘包,只读指定长度 byte[] bodyBytes = buffer.getBytes(17, 17 + header.getBodyLength()); // 获取序列化器 Serializer serializer = SerializerFactory.getSerializer( ProtocolMessageSerializerEnum.getEnumByKey(header.getSerializer()).getValue() ); // 反序列化 ProtocolMessageTypeEnum typeEnum = ProtocolMessageTypeEnum.getEnumByKey(header.getType()); switch (typeEnum) { case REQUEST: RpcRequest request = serializer.deserialize(bodyBytes, RpcRequest.class); return new ProtocolMessage<>(header, request); case RESPONSE: RpcResponse response = serializer.deserialize(bodyBytes, RpcResponse.class); return new ProtocolMessage<>(header, response); default: throw new RuntimeException("不支持的消息类型"); } }
八、TCP 粘包拆包解决方案
使用 Vert.x RecordParser 实现固定头+变长体解析。
原理
- 先读 17字节固定头
- 从头中获取 bodyLength
- 再读 bodyLength 字节
- 组合成完整包,再交给业务处理器
实现类:TcpBufferHandlerWrapper
▼java复制代码public class TcpBufferHandlerWrapper implements Handler<Buffer> { private final RecordParser recordParser; public TcpBufferHandlerWrapper(Handler<Buffer> bufferHandler) { recordParser = initRecordParser(bufferHandler); } private RecordParser initRecordParser(Handler<Buffer> bufferHandler) { // 先读固定长度头 RecordParser parser = RecordParser.newFixed(ProtocolConstant.MESSAGE_HEADER_LENGTH); parser.setHandler(buffer -> { // 第一次读取头 int bodyLength = buffer.getInt(13); // 切换读body parser.fixedSizeMode(bodyLength); // 继续读取... }); return parser; } @Override public void handle(Buffer buffer) { recordParser.handle(buffer); } }
九、服务端处理流程 TcpServerHandler
▼text复制代码接收 buffer → 解码 → 获取 RpcRequest → 反射调用 → 构造 RpcResponse → 编码发送
▼java复制代码public void handle(NetSocket netSocket) { TcpBufferHandlerWrapper wrapper = new TcpBufferHandlerWrapper(buffer -> { // 1. 解码 ProtocolMessage<RpcRequest> protocolMessage = (ProtocolMessage<RpcRequest>) ProtocolMessageDecoder.decode(buffer); RpcRequest request = protocolMessage.getBody(); // 2. 反射调用 Class<?> implClass = LocalRegistry.getService(request.getServiceName()); Method method = implClass.getMethod(request.getMethodName(), request.getParameterTypes()); Object result = method.invoke(implClass.newInstance(), request.getParameters()); // 3. 封装响应 RpcResponse response = new RpcResponse(); response.setData(result); // 4. 编码返回 ProtocolMessage.Header header = protocolMessage.getHeader(); header.setType((byte) ProtocolMessageTypeEnum.RESPONSE.getKey()); ProtocolMessage<RpcResponse> responseMsg = new ProtocolMessage<>(header, response); Buffer encodeBuffer = ProtocolMessageEncoder.encode(responseMsg); netSocket.write(encodeBuffer); }); netSocket.handler(wrapper); }
十、客户端请求流程 VertxTcpClient
▼text复制代码构造 RpcRequest → 构造 ProtocolMessage → 编码 → 发送 → 接收响应 → 解码 → 返回 RpcResponse
▼java复制代码public static RpcResponse doRequest(RpcRequest rpcRequest, ServiceMetaInfo serviceMetaInfo) { // 构造协议消息 ProtocolMessage<RpcRequest> protocolMessage = new ProtocolMessage<>(); ProtocolMessage.Header header = new ProtocolMessage.Header(); header.setMagic(ProtocolConstant.PROTOCOL_MAGIC); header.setSerializer(...); header.setType((byte) ProtocolMessageTypeEnum.REQUEST.getKey()); protocolMessage.setHeader(header); protocolMessage.setBody(rpcRequest); // 编码 Buffer buffer = ProtocolMessageEncoder.encode(protocolMessage); // 发送TCP请求 socket.write(buffer); // 接收响应 & 解码 ProtocolMessage<RpcResponse> responseMsg = (ProtocolMessage<RpcResponse>) ProtocolMessageDecoder.decode(buffer); return responseMsg.getBody(); }
评论
问答助学
相关内容
0个评论
全部评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
内容推荐
Day 29✅ 今天做了:Redis消息队列⏰ 明天计划:继续学消息队列
2
Day 17🧭行动:学习了Python最后一个知识点:异常处理🤓体会:Python 程序一旦发生异常,如果没有捕获处理,程序就会直接崩溃终止。使用"try-except"捕获异常,可以预先写好异常处理方案:比如打印友好提示、记录日志、释放资源,保证程序不会直接退出,还能继续运行。🧑💻代码:try:print("================================")# pri
2
Day1今天学习了java中if的使用
1
day 61今天来学校上课了,和教授谈了,可以接下RA工作,不过一周只能charge 10h, 工资很少,不过一个月600刀也勉强比没有好。今天的工作还没做完,最近在做知识库构造,Ui/UX上周就设计完,现在在做底座的建设今天也让Claude直接帮我把Neetcode150 题目+解题思路+代码给我整理了,我发现我真的只有找碎片化时间才能做题,我还不如先背先回忆(之前都做过)祝祖国母亲节日快乐,也
2
Day 1✅ 今天做了:⏰ 明天计划:📚 今日感悟:
0
