- 后端
- 消息队列
- 微服务
- 消息路由
【免费下载链接】CAP
Distributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern
导读
序列化是 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。几个关键实现细节值得注意:
- 序列化选项来自
CapOptions.JsonSerializerOptions:构造函数通过IOptions<CapOptions>注入,读取capOptions.Value.JsonSerializerOptions作为所有序列化/反序列化调用的选项对象。该属性在 CAP.Options.cs 中定义为public JsonSerializerOptions JsonSerializerOptions { get; } = new();,官方注释说明可以自定义它以控制 JSON 格式、命名策略、转换器(converter)等序列化行为; - 空载荷处理:
SerializeAsync中若message.Value == null,直接返回 Body 为空的TransportMessage(只携带 Headers);DeserializeAsync中若valueType == null或Body.Length == 0,同样返回Value == null的Message,避免对空体做无意义的 JSON 解析; - 字节级高效传输:
SerializeAsync使用JsonSerializer.SerializeToUtf8Bytes直接产出 UTF-8 字节数组,DeserializeAsync使用JsonSerializer.Deserialize(transportMessage.Body.Span, valueType, ...)从ReadOnlyMemory<byte>的 Span 上直接反序列化,贴合TransportMessage的字节承载设计; 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( /* ... */ );两点注册相关的实现细节值得说明:
TryAddSingleton的覆盖语义:AddCap内部对ISerializer使用的是TryAddSingleton<ISerializer, JsonUtf8Serializer>()(见 CAP.ServiceCollectionExtensions.cs),即仅在尚未注册时才注册默认实现。因此在上面的示例中,先注册自定义序列化器、再调用AddCap是推荐且干净的顺序——AddCap检测到ISerializer已有实现便不会覆盖;- 单例生命周期:
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生效。
注意事项与最佳实践
结合源码调用链,替换序列化器时有几点需要提前规划:
- 存储与传输同时受影响:序列化器既负责消息落库格式,也负责 Broker 字节流格式(发布端
MessageSender.SerializeAsync、消费端ConsumerRegister.DeserializeAsync均依赖同一ISerializer单例)。替换后,历史遗留消息若仍以旧格式(如 JSON)存储在数据库或队列中,新消费者可能无法正确还原,需要评估兼容与迁移方案; IsJsonType必须与Deserialize(object, Type)配套:订阅参数绑定先经IsJsonType判断再调用Deserialize,两者对“JSON 类型”的认定必须一致(默认实现统一以JsonElement为基准),否则订阅方法参数会落入TypeConverter/Convert.ChangeType兼容分支,可能导致非预期转换行为;- 异常消息存储同样走序列化器:在 IConsumerRegister.Default.cs 中,处理失败的消息会通过
_serializer.Serialize(message)序列化后调用StoreReceivedExceptionMessageAsync存入异常存储,因此自定义序列化器还需保证失败消息(含异常头)能够被正确序列化与反序列化; - 保持空载荷语义:默认实现允许
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
相关推荐
vLLM-Omni 扩散模型 CPU Offload 实战:从模型级到分布式层级级卸载的配置与源码解析
vLLM Omni 扩散模型 CPU Offload 实战:从模型级到分布式层级级卸载的配置与源码解析 本篇围绕 vLLM Omni 中扩散(Diffusion
后端消息队列微服务RestSharp 序列化指南:从 JSON/XML 默认序列化到自定义序列化器(.NET)
RestSharp 序列化指南:从 JSON/XML 默认序列化到自定义序列化器(.NET) RestSharp 作为 .NET 生态中最常用的 REST/HT
后端NoneBot2 消息处理机制详解:从消息序列到消息模板
NoneBot2 消息处理机制详解:从消息序列到消息模板 你是否曾在开发聊天机器人时遇到过这样的困扰:不同平台的消息格式五花八门,有的支持纯文本,有的支持富文本
后端即时通讯
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考