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 103✅ 今天做了:复习了多用户通信系统⏰ 明天计划:学习Java反射
0
Day 68时间19:00~ 22:00(3h)✅ 今天做了:Component注解、Mybatis配置、使用⏰ 明天计划:Lombok、Mapper映射、动态SQL📚 今日感悟:自动配置类DataSourceAutoConfiguration ,会读取properties文件,通过注解:@EnableConfigurationProperties(DataSourceProperties.cl
1
Day 19✅ 今天做了:MCP⏰ 明天计划:AI智能体构建📚 今日感悟:今天MCP问题有点多有点杂,明天找时间再捋一下。继续加油
0
Day 25✅ 今天做了:1、扇贝英语单词打卡2、英语听说读写、听力练习3、微信阅读15分钟4、编程导航学习⏰ 明天计划:待定📚 今日感悟:Keep going!
1
Day 104✅ 今天做了:学习了Java反射及快速入门⏰ 明天计划:继续学习Java反射
0
