You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用自定义序列化器解析KafkaFlow消息

实现KafkaFlow JSON的直接类型解析

目标JSON格式

{
    "actions": {
        "$type": "Full.ActionTypeName, supplying-assembly-name",
        "$values": [
            {
                "propname": "propvalue" //properties of ActionTypeName instance
            }
        ]
    },
    "resources": {
        "$type": "Full.ResourceTypeName, supplying-assembly-name",
        "$values": [
            {
                "propname": "propvalue" //properties of ResourceTypeName instance
            }
        ]
    }
}

现有C#定义与基础解析

实体类定义

using System.Text.Json;
using System.Text.Json.Serialization;

public class FlowMessage
{
    public FlowType Actions { get; set; }
    public FlowType Resources { get; set; }
}

public class FlowType
{
    [JsonPropertyName("$type")]
    public string TypeDescriptor { get; set; }
    [JsonPropertyName("$values")]
    public JsonElement[] Values { get; set; }
    public string TypeName() => TypeDescriptor.Split(",")[0];
    public Type Type() => System.Type.GetType(TypeName());
}

基础解析代码

var stream = File.Open("sample.json", FileMode.Open, FileAccess.Read);
JsonSerializerOptions seropt = new()
{
    PropertyNameCaseInsensitive = true,
    Converters = { new JsonStringEnumConverter() }
};

var foo = JsonSerializer.Deserialize(stream, typeof(FlowMessage), seropt);

核心需求

当前解析得到的是包含$type和$values节点的FlowType对象,需要让序列化器直接返回对应类型的对象数组,而非FlowType结构。

解决方案:自定义JsonConverter实现动态类型解析

通过自定义JsonConverterFactory,利用克隆JSON阅读器预读$type字段确定目标类型,再将$values数组解析为对应类型的实例数组。

步骤1:实现通用转换器工厂

public class FlowArrayConverterFactory : JsonConverterFactory
{
    public override bool CanConvert(Type typeToConvert)
    {
        return typeToConvert == typeof(object[]);
    }

    public override JsonConverter CreateConverter(Type typeToConvert, JsonSerializerOptions options)
    {
        return new FlowArrayConverter(options);
    }

    private class FlowArrayConverter : JsonConverter<object[]>
    {
        private readonly JsonSerializerOptions _options;

        public FlowArrayConverter(JsonSerializerOptions options)
        {
            // 克隆配置并移除自身,避免递归调用
            _options = new JsonSerializerOptions(options);
            _options.Converters.Remove(this);
        }

        public override object[] Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
        {
            // 克隆阅读器节点,预读$type字段
            using var doc = JsonDocument.ParseValue(ref reader);
            var root = doc.RootElement;

            // 解析目标类型
            if (!root.TryGetProperty("$type", out var typeElem))
                throw new JsonException("Missing $type property in flow array");
            var typeName = typeElem.GetString()!.Split(',')[0];
            var targetType = Type.GetType(typeName) ?? throw new JsonException($"Cannot resolve type {typeName}");

            // 解析$values为目标类型数组
            var valuesElem = root.GetProperty("$values");
            var targetArrayType = targetType.MakeArrayType();
            return (object[])JsonSerializer.Deserialize(valuesElem.GetRawText(), targetArrayType, _options)!;
        }

        public override void Write(Utf8JsonWriter writer, object[] value, JsonSerializerOptions options)
        {
            // 序列化反向逻辑:写入$type和$values
            if (value.Length == 0)
            {
                writer.WriteStartObject();
                writer.WriteEndObject();
                return;
            }

            var elemType = value[0].GetType();
            var typeDescriptor = $"{elemType.FullName}, {elemType.Assembly.GetName().Name}";

            writer.WriteStartObject();
            writer.WriteString("$type", typeDescriptor);
            writer.WritePropertyName("$values");
            JsonSerializer.Serialize(writer, value, elemType.MakeArrayType(), _options);
            writer.WriteEndObject();
        }
    }
}

步骤2:修改实体类并配置转换器

将FlowMessage的属性类型改为object[],并在序列化配置中添加自定义转换器:

public class FlowMessage
{
    public object[] Actions { get; set; }
    public object[] Resources { get; set; }
}

// 调整后的解析代码
var stream = File.Open("sample.json", FileMode.Open, FileAccess.Read);
JsonSerializerOptions seropt = new()
{
    PropertyNameCaseInsensitive = true,
    Converters = { new JsonStringEnumConverter(), new FlowArrayConverterFactory() }
};

var foo = JsonSerializer.Deserialize<FlowMessage>(stream, seropt);
// 直接强转为目标类型数组使用
var actions = (Full.ActionTypeName[])foo.Actions;

关键说明

  1. 克隆阅读器预读:通过JsonDocument.ParseValue克隆当前JSON节点,避免移动原阅读器指针,实现先读取$type确定目标类型,再解析数组内容。
  2. 转换器工厂:通过JsonConverterFactory动态适配任意目标类型,无需为每个类型单独编写转换器。
  3. 配置隔离:克隆序列化配置并移除自身转换器,防止递归调用导致死循环。

内容的提问来源于stack exchange,提问作者Peter Wone

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 16:27:48