☰
CAP 消息序列化机制详解:从默认 JSON 到自定义 ISerializer 扩展
2026/9/29 3:26:20 网站建设 项目流程
  • 后端
  • 消息队列
  • 微服务
  • 消息路由

【免费下载链接】CAP

Distributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern

项目地址:https://gitcode.com/gh_mirrors/ca/CAP
点击查看免费下载

导读

序列化是 CAP 分布式事务消息与事件总线的核心枢纽:它决定了消息内容以何种格式写入消息存储、以何种字节流送入消息队列,以及消费者侧如何把收到的字节还原成订阅方法参数。本文基于 CAP 官方用户指南中 Serialization 文档 展开,结合仓库源码梳理ISerializer接口的全貌、默认JsonUtf8Serializer的实现细节、消息在发布与消费两条链路中的序列化调用点,并给出一个完整的自定义序列化器实现与注册示例。读完后,你将能够根据业务需要替换 CAP 的消息序列化方案,并理解替换行为对整个消息管道的影响边界。

序列化在 CAP 中的角色

CAP 对外提供ISerializer接口统一承担消息的序列化与反序列化职责。官方文档明确指出:默认情况下,CAP 使用 JSON 对消息进行序列化,并将序列化结果存入数据库。这意味着无论你最终选用哪种消息队列(Kafka、RabbitMQ、Azure Service Bus、Redis Streams 等),消息落库与进出 Broker 的字节形态都由这一个接口控制,存储层与传输层无需关心具体的序列化格式。

CAP 的消息体系中有两个关键类型,理解它们的区别是掌握序列化接口的前提:

  • Message:应用层的消息抽象,由Headers(消息元数据字典,如 MessageId、MessageName、Group、CorrelationId 等)与Value(实际载荷对象)组成,是发布端发布、订阅方法接收的对象形态;
  • TransportMessage:传输层消息结构,是只读 struct,由Headers与Body(原始字节,通常是 UTF-8 编码的 JSON)组成,用于在消息管道中高效传递原始数据。

ISerializer的核心工作正是在这两种形态之间完成转换:发布时把Message(含对象载荷)变成携带字节 Body 的TransportMessage;消费时把TransportMessage的字节还原成订阅方法参数所需的强类型对象。

ISerializer 接口全貌

官方文档给出的自定义示例只展示了SerializeAsync与DeserializeAsync两个异步方法,但当前仓库中的 ISerializer.cs 接口实际包含 6 个成员,自定义实现时建议全部覆盖:

public interface ISerializer { /// <summary> /// 将 Message 序列化为字符串 /// </summary> string Serialize(Message message); /// <summary> /// 将 Message 序列化为 TransportMessage(异步) /// </summary> ValueTask<TransportMessage> SerializeAsync(Message message); /// <summary> /// 将字符串反序列化为 Message /// </summary> Message? Deserialize(string json); /// <summary> /// 将 TransportMessage 反序列化为 Message(异步) /// </summary> ValueTask<Message> DeserializeAsync(TransportMessage transportMessage, Type? valueType); /// <summary> /// 将指定对象按 valueType 反序列化 /// </summary> object? Deserialize(object value, Type valueType); /// <summary> /// 判断给定对象是否为 Json 类型(如 JToken 或 JsonElement,取决于实现的序列化器) /// </summary> bool IsJsonType(object jsonObject); }

各成员的用途可从源码调用点归纳如下:

接口成员主要调用场景
SerializeAsync(Message)发布端发送消息前,将业务对象序列化为传输字节
DeserializeAsync(TransportMessage, Type)消费端收到消息后,按订阅方法参数类型还原对象
Serialize(Message)消息处理失败、携带异常头时,将 Message 序列化为字符串存入异常存储
Deserialize(string)从存储读取字符串并还原为 Message
Deserialize(object, Type)订阅方法参数绑定:当参数值本身是 JsonElement 时按目标类型转换
IsJsonType(object)订阅调用器判断消息 Value 是否为 Json 类型,决定走反序列化还是类型转换分支

默认实现 JsonUtf8Serializer 的源码剖析

仓库在 CAP.ServiceCollectionExtensions.cs 中通过以下语句注册默认序列化器:

services.TryAddSingleton<ISerializer, JsonUtf8Serializer>();

JsonUtf8Serializer的实现位于 ISerializer.JsonUtf8.cs,底层基于System.Text.Json。几个关键实现细节值得注意:

  1. 序列化选项来自CapOptions.JsonSerializerOptions:构造函数通过IOptions<CapOptions>注入,读取capOptions.Value.JsonSerializerOptions作为所有序列化/反序列化调用的选项对象。该属性在 CAP.Options.cs 中定义为public JsonSerializerOptions JsonSerializerOptions { get; } = new();,官方注释说明可以自定义它以控制 JSON 格式、命名策略、转换器(converter)等序列化行为;
  2. 空载荷处理:SerializeAsync中若message.Value == null,直接返回 Body 为空的TransportMessage(只携带 Headers);DeserializeAsync中若valueType == null或Body.Length == 0,同样返回Value == null的Message,避免对空体做无意义的 JSON 解析;
  3. 字节级高效传输:SerializeAsync使用JsonSerializer.SerializeToUtf8Bytes直接产出 UTF-8 字节数组,DeserializeAsync使用JsonSerializer.Deserialize(transportMessage.Body.Span, valueType, ...)从ReadOnlyMemory<byte>的 Span 上直接反序列化,贴合TransportMessage的字节承载设计;
  4. IsJsonType的实现:默认实现返回jsonObject is JsonElement,即识别System.Text.Json的JsonElement为 JSON 类型。

发布链路的调用点

在 IMessageSender.Default.cs 中,SendWithoutRetryAsync发送消息的第一步就是:

var transportMsg = await _serializer.SerializeAsync(message.Origin).ConfigureAwait(false);

也就是说,ICapPublisher.PublishAsync之后,消息先进入存储(落库时已按序列化格式存储),再由MessageSender通过ISerializer.SerializeAsync把Message转成TransportMessage交给ITransport.SendAsync送入消息队列。这里使用的是容器注入的ISerializer单例(构造函数经serviceProvider.GetRequiredService<ISerializer>()解析,见 IMessageSender.Default.cs),因此替换序列化器会同时影响存储与传输两个环节。

消费链路的调用点

在 IConsumerRegister.Default.cs 中,消费端取出订阅方法描述后,按第一个非 CAP 内置参数的参数类型进行反序列化:

var type = executor!.Parameters.FirstOrDefault(x => x.IsFromCap == false)?.ParameterType; message = await _serializer.DeserializeAsync(transportMessage, type);

随后在 ISubscribeInvoker.Default.cs 的订阅方法参数绑定阶段,调用器会先通过_serializer.IsJsonType(message.Value)判断消息载荷是否为 JSON 类型:

  • 是 JSON 类型:调用_serializer.Deserialize(message.Value, parameterDescriptor.ParameterType)按目标参数类型转换;
  • 不是 JSON 类型:走TypeDescriptor.GetConverter、IsInstanceOfType、Convert.ChangeType等兼容转换分支。

因此IsJsonType的实现必须与Deserialize(object, Type)保持语义一致——前者识别出的 JSON 对象类型(如JsonElement)正是后者能够直接处理的类型。默认的JsonUtf8Serializer对此已给出标准实现模板,接口注释中也附带了System.Text.Json场景的示例代码。

自定义序列化:完整实现示例

官方文档给出了自定义序列化器的骨架,下面基于接口全貌补全为一个可直接编译、可直接注册的完整实现(这里以System.Text.Json风格为例;若改用Newtonsoft.Json,IsJsonType通常应判断JToken):

public class YourSerializer : ISerializer { // 可以注入自定义选项,例如统一的 JsonSerializerSettings public YourSerializer() { } public ValueTask<TransportMessage> SerializeAsync(Message message) { if (message == null) throw new ArgumentNullException(nameof(message)); // 空载荷:仅携带 Headers,Body 为空 if (message.Value == null) return new ValueTask<TransportMessage>(new TransportMessage(message.Headers, null)); // 把业务对象序列化为 UTF-8 字节,作为 TransportMessage.Body var jsonBytes = JsonSerializer.SerializeToUtf8Bytes(message.Value); return new ValueTask<TransportMessage>(new TransportMessage(message.Headers, jsonBytes)); } public ValueTask<Message> DeserializeAsync(TransportMessage transportMessage, Type? valueType) { if (valueType == null || transportMessage.Body.Length == 0) return new ValueTask<Message>(new Message(transportMessage.Headers, null)); var obj = JsonSerializer.Deserialize(transportMessage.Body.Span, valueType); return new ValueTask<Message>(new Message(transportMessage.Headers, obj)); } public string Serialize(Message message) { return JsonSerializer.Serialize(message); } public Message? Deserialize(string json) { return JsonSerializer.Deserialize<Message>(json); } public object? Deserialize(object value, Type valueType) { if (value is JsonElement jsonElement) return jsonElement.Deserialize(valueType); throw new NotSupportedException("Type is not of type JsonElement"); } public bool IsJsonType(object jsonObject) { return jsonObject is JsonElement; } }

注册自定义序列化器

按照官方文档,将自定义实现注册到依赖注入容器,然后再调用AddCap:

services.AddSingleton<ISerializer, YourSerializer>(); services.AddCap( /* ... */ );

两点注册相关的实现细节值得说明:

  1. TryAddSingleton的覆盖语义:AddCap内部对ISerializer使用的是TryAddSingleton<ISerializer, JsonUtf8Serializer>()(见 CAP.ServiceCollectionExtensions.cs),即仅在尚未注册时才注册默认实现。因此在上面的示例中,先注册自定义序列化器、再调用AddCap是推荐且干净的顺序——AddCap检测到ISerializer已有实现便不会覆盖;
  2. 单例生命周期:ISerializer在 CAP 内部以单例方式解析(MessageSender、SubscribeInvoker、ConsumerRegister均通过GetRequiredService<ISerializer>()获取),注册时也应使用AddSingleton,保证发布与消费链路复用同一序列化实例,避免重复创建带来的状态不一致。

自定义 JSON 序列化选项(不改实现)

如果不需要替换整体序列化方案,只是希望调整默认 JSON 行为的细节(如命名策略、时间格式、忽略空值、追加自定义 Converter),可以直接配置CapOptions.JsonSerializerOptions:

services.AddCap(options => { // 示例:统一使用 camelCase 命名策略,并添加自定义转换器 options.JsonSerializerOptions.PropertyNamingPolicy = JsonNamingPolicy.CamelCase; options.JsonSerializerOptions.Converters.Add(new MyCustomConverter()); // 其余 CAP 配置 ... });

该选项对象会被注入JsonUtf8Serializer构造函数的IOptions<CapOptions>中,作用于所有JsonSerializer.SerializeToUtf8Bytes/Deserialize调用(见 ISerializer.JsonUtf8.cs)。注意该属性为只读初始化(get;私有 set),只能在AddCap配置回调中修改,且仅对默认的JsonUtf8Serializer生效。

注意事项与最佳实践

结合源码调用链,替换序列化器时有几点需要提前规划:

  1. 存储与传输同时受影响:序列化器既负责消息落库格式,也负责 Broker 字节流格式(发布端MessageSender.SerializeAsync、消费端ConsumerRegister.DeserializeAsync均依赖同一ISerializer单例)。替换后,历史遗留消息若仍以旧格式(如 JSON)存储在数据库或队列中,新消费者可能无法正确还原,需要评估兼容与迁移方案;
  2. IsJsonType必须与Deserialize(object, Type)配套:订阅参数绑定先经IsJsonType判断再调用Deserialize,两者对“JSON 类型”的认定必须一致(默认实现统一以JsonElement为基准),否则订阅方法参数会落入TypeConverter/Convert.ChangeType兼容分支,可能导致非预期转换行为;
  3. 异常消息存储同样走序列化器:在 IConsumerRegister.Default.cs 中,处理失败的消息会通过_serializer.Serialize(message)序列化后调用StoreReceivedExceptionMessageAsync存入异常存储,因此自定义序列化器还需保证失败消息(含异常头)能够被正确序列化与反序列化;
  4. 保持空载荷语义:默认实现允许Value == null/Body为空的消息存在(用于仅携带 Headers 的场景),自定义实现应保持这一语义,避免对空体强行解析而抛异常。

小结

序列化是 CAP 消息管道中贯穿“应用对象 → 存储/传输字节 → 订阅参数”的关键抽象。默认的JsonUtf8Serializer基于System.Text.Json实现,其选项可通过CapOptions.JsonSerializerOptions调整;需要整体替换格式时,实现ISerializer的全部成员并用AddSingleton在AddCap之前注册即可。理解Message与TransportMessage两种形态的差异,以及SerializeAsync/DeserializeAsync/IsJsonType在发布、消费、参数绑定三条路径上的调用位置,是安全实施自定义序列化方案的前提。相关完整实现与测试可继续参阅 ISerializer.JsonUtf8.cs、ISerializer.cs 以及 ISubscribeInvoker.Default.cs。

  • 后端
  • 消息队列
  • 微服务
  • 消息路由

【免费下载链接】CAP

Distributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern

项目地址:https://gitcode.com/gh_mirrors/ca/CAP
点击查看免费下载

相关推荐

上一篇:NeteaseCloudMusicFlac无损音乐批量下载教程:一张网易云歌单8分钟下完全部FLAC
下一篇:Colima 自动化脚本实战:为 bootstrap、CI 与部署流程编写非交互式驱动

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询