如何从Avsc文件生成可被SchemaRegistry.Serdes.Avro的AvroSerializer处理的C#类
从Avsc文件生成兼容SchemaRegistry.Serdes.Avro的C#类
一、使用官方代码生成工具
Confluent提供的avrogen命令行工具是生成符合要求C#类的首选方式,步骤如下:
- 安装工具:通过dotnet全局工具安装
dotnet tool install -g Confluent.Apache.Avro.Tools - 执行生成命令:
参数说明: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
相关产品推荐
相关产品推荐

