系列

第03篇:手写JavaRPC框架之搞定序列化

2026-01-295 min read手写RPC框架
Mojito

天下代码一大抄, 抄来抄去有提高, 看你会抄不会抄!

一、前言

天下代码一大抄, 抄来抄去有提高, 看你会抄不会抄!从本篇开始后面的所有章节都是实战环节,每节一个小目标,最终我们实现完整的JavaRPC的框架,然后发布maven仓库,感兴趣的同学可以下载研究。大家如果想要获取源码的话可以私信: RPC,自动回复仓库地址。

其实这些东西并没有什么难度,只要从头到尾跟着我们一起coding,其实就会发现不过如此。所以就算是新手也不要有心里负担。还是那句话: "天下代码一大抄, 抄来抄去有提高, 看你会抄不会抄"。主要的是思想,而不是死记硬背。所以我们来主要来学习设计思想,具体的代码,收藏用的时候看下就好了。

下面我们废话就不多少,直接开始吧。

二、目标

2.1 目标介绍

RPC框架中最基础的一个功能就是通讯,而通讯的本质就是发送方将数据转换成二进制的数据流,然后通过网络管道将数据发送给服务方,服务方在通过读取管道中的二进制数据最终将数据从二进制转换成Java对象,然后供Java系统处理,处理完成后再将Java结果对象转换成二进制数据通过管道发送回去。

这里有两个重点,第一个是网络通信,就是图中的网路连接管道,第二个是管道中的二进制数据。本篇文章我们先研究后者,就是通过代码将Java对象转换成二进制数据。因为通信比较难,内容也较多,所以我们先易后难,难的放到下一篇文章在说。

这里有两个术语: 序列化和反序列化

如下实例User类,转换成二进制数据就是一个数组。这个Java对象转换二进制数据的过程叫做序列化。 而二进制数据转换成Java对象的过程叫做反序列化。

2.2 两个小目标

  • 第一个目标就是我们实现多种序列化的能力。
  • 第二个目标是来选择一个最优的序列化方案。

为什么说要选择最优的序列化方案呢? 因为我们要适配下面这两种场景。

  1. 第一种是对性能要有比较高的,这种情况就要求我们序列化和反序列的速度要足够快
  2. 第二种就是对内存和空间要有比较高的,就要求我们的序列化后的数据要足够的小。

下面就开始Coding了。

三、设计

3.1 工程结构

首先我们根据上一篇文章中定义的项目分层结构,先把需要的所有层给创建出来。然后实现序列化。

text
1mojito/mojito-net/src/main/java/cn/lxchinesszz/mojito on  master [!+?] 2. 3├── api 4├── business 5├── codec 6├── exception 7│   ├── BusinessServerHandlerException.java 8│   ├── DeserializeException.java 9│   ├── HttpsTokenFileException.java 10│   ├── ProtocolException.java 11│   ├── RemotingException.java 12│   ├── SerializeException.java 13│   └── SignatureException.java 14├── exchange 15└── serialize 16 ├── AbstractSerializer.java 17 ├── Serializer.java 18 └── impl 19 ├── Hession2ObjectSerializer.java 20 ├── HessionObjectSerializer.java 21 ├── NettyCompactObjectSerializer.java 22 ├── NettyObjectSerializer.java 23 └── ProtostuffObjectSerializer.java 24

3.2 代码结构

首先我们先定义接口,为什么定义接口呢? 方便后面的扩展。 接口也非常的简单就三个方法。

  1. 负责将任意对象转换成二进制数组
  2. 将二进制数组转成Object对象
  3. 将二进制数组转换成指定的Java对象

UML图设计

3.3 实现

  • ProtostuffObjectSerializer 使用谷歌开源的序列化库,特点占用极小。
java
1public class ProtostuffObjectSerializer extends AbstractSerializer { 2 3 /** 4 * 线程数会有限制,不会无穷大的使用 5 */ 6 private static final ThreadLocal<LinkedBuffer> BUFFER = InheritableThreadLocal.withInitial(() -> 7 LinkedBuffer.allocate(LinkedBuffer.DEFAULT_BUFFER_SIZE)); 8 9@Override 10 @SuppressWarnings("unchecked") 11 public byte[] doSerialize(Object dataObject) throws SerializeException { 12 // // RuntimeSchema 懒加载内置缓存,所以我们不用在缓存了 13 Schema<Object> schema = (Schema<Object>) RuntimeSchema.getSchema(dataObject.getClass()); 14LinkedBuffer linkedBuffer = BUFFER.get(); 15byte[] bytes = ProtostuffIOUtil.toByteArray(dataObject, schema, linkedBuffer); 16linkedBuffer.clear(); 17return bytes; 18} 19 20 @Override 21 public Object doDeserialize(byte[] data) throws DeserializeException { 22 throw new UnsupportedOperationException(getClass() + "必须指定反序列化类型"); 23} 24 25 @Override 26 public <T> T deserialize(byte[] data, Class<T> dataType) throws DeserializeException { 27 return schema(data, dataType); 28} 29 30 public <T> T schema(byte[] data, Class<T> dataType) throws DeserializeException { 31 try { 32 Schema<T> schema = RuntimeSchema.getSchema(dataType); 33T dataTypeObj = dataType.newInstance(); 34ProtostuffIOUtil.mergeFrom(data, dataTypeObj, schema); 35return dataTypeObj; 36} catch (InstantiationException | IllegalAccessException e) { 37 throw new DeserializeException(e); 38} 39 } 40}
  • NettyCompactObjectSerializer Netty原生支持的序列化协议
java
1public class NettyCompactObjectSerializer extends AbstractSerializer { 2 3 @Override 4 public byte[] doSerialize(Object dataObject) throws SerializeException { 5 try (ByteArrayOutputStream dataArr = new ByteArrayOutputStream(); 6 CompactObjectOutputStream oeo = new CompactObjectOutputStream(dataArr)) { 7 oeo.writeObject(dataObject); 8oeo.flush(); 9return dataArr.toByteArray(); 10} catch (IOException e) { 11 throw new SerializeException(e); 12} 13 } 14 15 @Override 16 public Object doDeserialize(byte[] data) throws DeserializeException { 17 Object o; 18try (CompactObjectInputStream odi = new CompactObjectInputStream(new ByteArrayInputStream(data), ClassResolvers.cacheDisabled(null))) { 19 o = odi.readObject(); 20} catch (ClassNotFoundException | IOException e) { 21 throw new DeserializeException(e); 22} 23 return o; 24} 25 26 private static class CompactObjectOutputStream extends ObjectOutputStream { 27 28 static final int TYPE_FAT_DESCRIPTOR = 0; 29static final int TYPE_THIN_DESCRIPTOR = 1; 30 31CompactObjectOutputStream(OutputStream out) throws IOException { 32 super(out); 33} 34 35 @Override 36 protected void writeStreamHeader() throws IOException { 37 writeByte(STREAM_VERSION); 38} 39 40 @Override 41 protected void writeClassDescriptor(ObjectStreamClass desc) throws IOException { 42 Class<?> clazz = desc.forClass(); 43if (clazz.isPrimitive() || clazz.isArray() || clazz.isInterface() || 44 desc.getSerialVersionUID() == 0) { 45 write(TYPE_FAT_DESCRIPTOR); 46super.writeClassDescriptor(desc); 47} else { 48 write(TYPE_THIN_DESCRIPTOR); 49writeUTF(desc.getName()); 50} 51 } 52 } 53 54 /** 55 * 压缩 56 */ 57 private static class CompactObjectInputStream extends ObjectInputStream { 58 59 static final int TYPE_FAT_DESCRIPTOR = 0; 60 61static final int TYPE_THIN_DESCRIPTOR = 1; 62 63private final ClassResolver classResolver; 64 65CompactObjectInputStream(InputStream in, ClassResolver classResolver) throws IOException { 66 super(in); 67this.classResolver = classResolver; 68} 69 70 @Override 71 protected void readStreamHeader() throws IOException { 72 int version = readByte() & 0xFF; 73if (version != STREAM_VERSION) { 74 throw new StreamCorruptedException( 75 "Unsupported version: " + version); 76} 77 } 78 79 @Override 80 protected ObjectStreamClass readClassDescriptor() 81 throws IOException, ClassNotFoundException { 82 int type = read(); 83if (type < 0) { 84 throw new EOFException(); 85} 86 switch (type) { 87 case TYPE_FAT_DESCRIPTOR: 88 return super.readClassDescriptor(); 89case TYPE_THIN_DESCRIPTOR: 90 String className = readUTF(); 91Class<?> clazz = classResolver.resolve(className); 92return ObjectStreamClass.lookupAny(clazz); 93default: 94 throw new StreamCorruptedException( 95 "Unexpected class descriptor type: " + type); 96} 97 } 98 99 @Override 100 protected Class<?> resolveClass(ObjectStreamClass desc) throws IOException, ClassNotFoundException { 101 Class<?> clazz; 102try { 103 clazz = classResolver.resolve(desc.getName()); 104} catch (ClassNotFoundException ignored) { 105 clazz = super.resolveClass(desc); 106} 107 108 return clazz; 109} 110 } 111} 112
  • HessionObjectSerializer 基于Hession序列化协议的封装
java
1public class HessionObjectSerializer extends AbstractSerializer { 2 3 @Override 4 public byte[] doSerialize(Object dataObject) throws SerializeException { 5 HessianOutput oeo = null; 6try (ByteArrayOutputStream dataArr = new ByteArrayOutputStream()) { 7 oeo = new HessianOutput(dataArr); 8oeo.writeObject(dataObject); 9oeo.flush(); 10return dataArr.toByteArray(); 11} catch (IOException e) { 12 throw new SerializeException(e); 13} finally { 14 if (oeo != null) { 15 try { 16 oeo.close(); 17} catch (IOException e) { 18 throw new SerializeException(e); 19} 20 } 21 } 22 } 23 24 @Override 25 public Object doDeserialize(byte[] data) throws DeserializeException { 26 Object o; 27HessianInput odi = new HessianInput(new ByteArrayInputStream(data)); 28try { 29 o = odi.readObject(); 30} catch (IOException e) { 31 throw new DeserializeException(e); 32} 33 return o; 34} 35 36} 37
  • Hession2ObjectSerializer 对Hession2序列化的封装
java
1public class Hession2ObjectSerializer extends AbstractSerializer { 2 3 @Override 4 public byte[] doSerialize(Object dataObject) throws SerializeException { 5 ByteArrayOutputStream dataArr = new ByteArrayOutputStream(); 6Hessian2Output oeo = null; 7try { 8 oeo = new Hessian2Output(dataArr); 9oeo.writeObject(dataObject); 10oeo.flush(); 11} catch (IOException e) { 12 throw new SerializeException(e); 13} finally { 14 if (oeo != null) { 15 try { 16 oeo.close(); 17} catch (IOException e) { 18 throw new SerializeException(e); 19} 20 } 21 } 22 return dataArr.toByteArray(); 23} 24 25 @Override 26 public Object doDeserialize(byte[] data) throws DeserializeException { 27 Object o; 28Hessian2Input odi = new Hessian2Input(new ByteArrayInputStream(data)); 29try { 30 o = odi.readObject(); 31} catch (IOException e) { 32 throw new DeserializeException(e); 33} 34 return o; 35} 36 37}

如果有优化的建议,希望多多评论,指点一二。以上代码已上传,代码较多就不这里展示。感兴趣的同学可以访问下面链接,查看代码。

仓库实现

四、性能分析

4.1 性能分析

那么以上这几种序列化协议那种速度最好呢?

首先我们先看序列化的速度对比。

使用20个线程,预热3秒,执行5秒对比下吞吐量及内存使用。这里我们只截图排名第一的实现。具体的对比数据已放在github上面,感兴趣的可以自己拉下来看下。

  • ProtostuffObjectSerializer 吞吐量 10,774,518 / s
  • Hession2ObjectSerializer 吞吐量 3,200,559 / s
  • HessionObjectSerializer 吞吐量 3,307,595 / s
  • NettyCompactObjectSerializer 吞吐量 874,793 / s

从上面中的数据我们能看出Protostuff是遥遥领先,这得益于提前把数据类型,提前进行了处理并缓存,其次是友好的API涉及,可以让我们避免频繁生成堆空间。

更多详细的数据对比,可以点下面链接查看。

详细数据

4.2 序列化长度

前面Protostuff在性能上是遥遥领先,那么在数据压缩比例上究竟谁能胜出呢? 相对于前面性能测试,这种大小测试是比较容易测试, 我们直接执行下面代码,查看结果。

java
1 @Test 2 @DisplayName("序列化数据大小对比") 3 public void test7() { 4 ProtostuffObjectSerializer nettyObjectSerializer = new ProtostuffObjectSerializer(); 5byte[] jay = nettyObjectSerializer.serialize(new UserSerializable("周杰伦", 42)); 6ColorConsole.colorPrintln("{},序列化数据长度:{}", "Protostuff", jay.length); 7 8ColorConsole.colorPrintln("{},序列化数据长度:{}", "Hession2", 9 new Hession2ObjectSerializer().serialize(new UserSerializable("周杰伦", 42)).length); 10 11ColorConsole.colorPrintln("{},序列化数据长度:{}", "Hession", 12 new HessionObjectSerializer().serialize(new UserSerializable("周杰伦", 42)).length); 13 14ColorConsole.colorPrintln("{},序列化数据长度:{}", "NettyCompact", 15 new NettyCompactObjectSerializer().serialize(new UserSerializable("周杰伦", 42)).length); 16}
实现类长度
ProtostuffObjectSerializer13字节
Hession2ObjectSerializer88字节
HessionObjectSerializer98字节
NettyCompactObjectSerializer132字节

还是Protostuff大幅度领先,那么我们的默认序列化协议就使用了吧。是不是很简单呢?

五、总结

通过这一节的实现,我们已经可以使用序列化工具,将Java对象转换成二进制数据了,那么下一篇我们就来实现如果创建网络连接管道吧。

那么你准备好跟我一起Coding了吗?