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

如何从Avsc文件生成可被SchemaRegistry.Serdes.Avro的AvroSerializer处理的C#类

从Avsc文件生成兼容SchemaRegistry.Serdes.Avro的C#类

一、使用官方代码生成工具

Confluent提供的avrogen命令行工具是生成符合要求C#类的首选方式,步骤如下:

  1. 安装工具:通过dotnet全局工具安装
    dotnet tool install -g Confluent.Apache.Avro.Tools
    
  2. 执行生成命令:
    avrogen -s 你的schema文件.avsc 输出目录路径
    
    参数说明:-s指定Avsc文件路径,后面的路径是生成类的存放目录。

二、生成类的核心特征(匹配示例结构)

生成的类会自动满足AvroSerializer的处理要求,核心规则包括:

  • 实现ISpecificRecord接口
  • 每个属性标记[AvroField("对应Avro字段名")]特性,保证字段映射准确
  • 实现ISpecificRecord的两个方法:
    • object Get(int fieldPos):根据字段位置返回对应属性值
    • void Put(int fieldPos, object value):根据字段位置设置对应属性值
  • 包含默认构造函数,可选全参构造函数(方便实例化)

举个对应示例:
假设Avsc定义如下:

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "Id", "type": "int"},
    {"name": "Name", "type": "string"}
  ]
}

生成的C#类结构大致为:

using System;
using Avro;

public class User : ISpecificRecord
{
    public static Schema _SCHEMA = Schema.Parse(@"{""type"":""record"",""name"":""User"",""fields"":[{""name"":""Id"",""type"":""int""},{""name"":""Name"",""type"":""string""}]}");

    [AvroField("Id")]
    public int Id { get; set; }

    [AvroField("Name")]
    public string Name { get; set; }

    public Schema Schema => _SCHEMA;

    public object Get(int fieldPos)
    {
        return fieldPos switch
        {
            0 => Id,
            1 => Name,
            _ => throw new ArgumentOutOfRangeException(nameof(fieldPos))
        };
    }

    public void Put(int fieldPos, object value)
    {
        switch (fieldPos)
        {
            case 0: Id = (int)value; break;
            case 1: Name = (string)value; break;
            default: throw new ArgumentOutOfRangeException(nameof(fieldPos));
        }
    }
}

三、配合AvroSerializer使用生成类

在Kafka生产者/消费者中,直接将生成的类作为泛型参数传入即可:

  • 生产者示例:
    var producerConfig = new ProducerConfig { BootstrapServers = "localhost:9092" };
    var schemaRegistryConfig = new SchemaRegistryConfig { Url = "http://localhost:8081" };
    
    using var schemaRegistry = new CachedSchemaRegistryClient(schemaRegistryConfig);
    using var producer = new ProducerBuilder<string, User>(producerConfig)
        .SetValueSerializer(new AvroSerializer<User>(schemaRegistry))
        .Build();
    
    await producer.ProduceAsync("user-topic", new Message<string, User> { Key = "1", Value = new User { Id = 1, Name = "Alice" } });
    
  • 消费者示例:
    var consumerConfig = new ConsumerConfig { BootstrapServers = "localhost:9092", GroupId = "user-group" };
    var schemaRegistryConfig = new SchemaRegistryConfig { Url = "http://localhost:8081" };
    
    using var schemaRegistry = new CachedSchemaRegistryClient(schemaRegistryConfig);
    using var consumer = new ConsumerBuilder<string, User>(consumerConfig)
        .SetValueDeserializer(new AvroDeserializer<User>(schemaRegistry).AsSyncOverAsync())
        .Build();
    
    consumer.Subscribe("user-topic");
    var consumeResult = consumer.Consume();
    User user = consumeResult.Message.Value;
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 06:29:52