myrpc 学习笔记-006 自定义协议

项目地址 欢迎访问

https://gitee.com/longlong5/myrpc

笔记总览 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

  • 自定义协议
  • 实现协议消息编解码器

v6.0.0.png

一、为什么要自定义 RPC 协议?

  1. HTTP 协议头太重,纯 RPC 不需要那么多信息
  2. 自定义协议更轻量、更快、更安全
  3. 方便扩展序列化器、消息类型、状态码
  4. 解决 TCP 粘包/拆包问题
  5. 统一消息格式,方便客户端 + 服务端通信

二、自定义协议整体结构

协议采用 固定头 + 可变体 格式:

text
复制代码
[ 消息头(17字节 固定)] + [ 消息体(长度不固定)]

1. 消息头结构(17字节)

偏移字段长度作用
0magic(魔数)1byte验证消息合法性
1version(版本)1byte协议升级
2serializer(序列化器)1bytejdk/json/kryo/hessian
3type(消息类型)1byte请求/响应/心跳
4status(状态)1byte成功/失败
5~12requestId8byte请求唯一ID
13~16bodyLength4byte消息体长度

总长度: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

流程

  1. 写入固定头(magic、version、serializer、type、status、requestId)
  2. 获取序列化器(根据头信息)
  3. 序列化消息体 → bodyBytes
  4. 写入bodyLength
  5. 写入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

功能

BufferProtocolMessage

流程

  1. 读取固定头
  2. 校验魔数
  3. 获取序列化器消息类型
  4. 读取bodyLength
  5. 截取bodyBytes
  6. 反序列化为 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 实现固定头+变长体解析。

原理

  1. 先读 17字节固定头
  2. 从头中获取 bodyLength
  3. 再读 bodyLength 字节
  4. 组合成完整包,再交给业务处理器

实现类: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个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
longlong
下载 APP