ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

CAP 消息序列化机制详解:从默认 JSON 到自定义 ISerializer 扩展

2026/9/29 3:26:31 拓冰建站 浏览量
CAP 消息序列化机制详解:从默认 JSON 到自定义 ISerializer 扩展 后端消息队列微服务消息路由【免费下载链接】CAPDistributed 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 ValueTaskTransportMessage SerializeAsync(Message message); /// summary /// 将字符串反序列化为 Message /// /summary Message? Deserialize(string json); /// summary /// 将 TransportMessage 反序列化为 Message异步 /// /summary ValueTaskMessage 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)从存储读取字符串并还原为 MessageDeserialize(object, Type)订阅方法参数绑定当参数值本身是 JsonElement 时按目标类型转换IsJsonType(object)订阅调用器判断消息 Value 是否为 Json 类型决定走反序列化还是类型转换分支默认实现 JsonUtf8Serializer 的源码剖析仓库在 CAP.ServiceCollectionExtensions.cs 中通过以下语句注册默认序列化器services.TryAddSingletonISerializer, JsonUtf8Serializer();JsonUtf8Serializer的实现位于 ISerializer.JsonUtf8.cs底层基于System.Text.Json。几个关键实现细节值得注意序列化选项来自CapOptions.JsonSerializerOptions构造函数通过IOptionsCapOptions注入读取capOptions.Value.JsonSerializerOptions作为所有序列化/反序列化调用的选项对象。该属性在 CAP.Options.cs 中定义为public JsonSerializerOptions JsonSerializerOptions { get; } new();官方注释说明可以自定义它以控制 JSON 格式、命名策略、转换器converter等序列化行为空载荷处理SerializeAsync中若message.Value null直接返回 Body 为空的TransportMessage只携带 HeadersDeserializeAsync中若valueType null或Body.Length 0同样返回Value null的Message避免对空体做无意义的 JSON 解析字节级高效传输SerializeAsync使用JsonSerializer.SerializeToUtf8Bytes直接产出 UTF-8 字节数组DeserializeAsync使用JsonSerializer.Deserialize(transportMessage.Body.Span, valueType, ...)从ReadOnlyMemorybyte的 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.GetRequiredServiceISerializer()解析见 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.JsonIsJsonType通常应判断JTokenpublic class YourSerializer : ISerializer { // 可以注入自定义选项例如统一的 JsonSerializerSettings public YourSerializer() { } public ValueTaskTransportMessage SerializeAsync(Message message) { if (message null) throw new ArgumentNullException(nameof(message)); // 空载荷仅携带 HeadersBody 为空 if (message.Value null) return new ValueTaskTransportMessage(new TransportMessage(message.Headers, null)); // 把业务对象序列化为 UTF-8 字节作为 TransportMessage.Body var jsonBytes JsonSerializer.SerializeToUtf8Bytes(message.Value); return new ValueTaskTransportMessage(new TransportMessage(message.Headers, jsonBytes)); } public ValueTaskMessage DeserializeAsync(TransportMessage transportMessage, Type? valueType) { if (valueType null || transportMessage.Body.Length 0) return new ValueTaskMessage(new Message(transportMessage.Headers, null)); var obj JsonSerializer.Deserialize(transportMessage.Body.Span, valueType); return new ValueTaskMessage(new Message(transportMessage.Headers, obj)); } public string Serialize(Message message) { return JsonSerializer.Serialize(message); } public Message? Deserialize(string json) { return JsonSerializer.DeserializeMessage(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; } }注册自定义序列化器按照官方文档将自定义实现注册到依赖注入容器然后再调用AddCapservices.AddSingletonISerializer, YourSerializer(); services.AddCap( /* ... */ );两点注册相关的实现细节值得说明TryAddSingleton的覆盖语义AddCap内部对ISerializer使用的是TryAddSingletonISerializer, JsonUtf8Serializer()见 CAP.ServiceCollectionExtensions.cs即仅在尚未注册时才注册默认实现。因此在上面的示例中先注册自定义序列化器、再调用AddCap是推荐且干净的顺序——AddCap检测到ISerializer已有实现便不会覆盖单例生命周期ISerializer在 CAP 内部以单例方式解析MessageSender、SubscribeInvoker、ConsumerRegister均通过GetRequiredServiceISerializer()获取注册时也应使用AddSingleton保证发布与消费链路复用同一序列化实例避免重复创建带来的状态不一致。自定义 JSON 序列化选项不改实现如果不需要替换整体序列化方案只是希望调整默认 JSON 行为的细节如命名策略、时间格式、忽略空值、追加自定义 Converter可以直接配置CapOptions.JsonSerializerOptionsservices.AddCap(options { // 示例统一使用 camelCase 命名策略并添加自定义转换器 options.JsonSerializerOptions.PropertyNamingPolicy JsonNamingPolicy.CamelCase; options.JsonSerializerOptions.Converters.Add(new MyCustomConverter()); // 其余 CAP 配置 ... });该选项对象会被注入JsonUtf8Serializer构造函数的IOptionsCapOptions中作用于所有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。赞分享后端消息队列微服务消息路由【免费下载链接】CAPDistributed transaction solution in micro-service base on eventually consistency, also an eventbus with Outbox pattern项目地址https://gitcode.com/gh_mirrors/ca/CAP点击查看免费下载相关推荐vLLM-Omni 扩散模型 CPU Offload 实战从模型级到分布式层级级卸载的配置与源码解析vLLM Omni 扩散模型 CPU Offload 实战从模型级到分布式层级级卸载的配置与源码解析 本篇围绕 vLLM Omni 中扩散Diffusion后端消息队列微服务RestSharp 序列化指南从 JSON/XML 默认序列化到自定义序列化器.NETRestSharp 序列化指南从 JSON/XML 默认序列化到自定义序列化器.NET RestSharp 作为 .NET 生态中最常用的 REST/HT后端NoneBot2 消息处理机制详解从消息序列到消息模板NoneBot2 消息处理机制详解从消息序列到消息模板 你是否曾在开发聊天机器人时遇到过这样的困扰不同平台的消息格式五花八门有的支持纯文本有的支持富文本后端即时通讯上一篇NeteaseCloudMusicFlac无损音乐批量下载教程一张网易云歌单8分钟下完全部FLAC下一篇Colima 自动化脚本实战为 bootstrap、CI 与部署流程编写非交互式驱动创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考