.Net集成Kafka时Avro序列化报错:Local Value serialization error
问题描述
在.NET环境中使用Kafka,已配置带证书认证的Schema Registry和Producer,尝试发送MLFixtureEmp类型消息到Kafka Topic时出现序列化错误。
配置及生产代码
var schemaRegistryConfig = new SchemaRegistryConfig { Url = "", SslKeystoreLocation = Path.Combine(Directory.GetCurrentDirectory(), @"Certs/decodedca.p12"), SslKeystorePassword = "", EnableSslCertificateVerification = false }; var producerConfig = new ProducerConfig { BootstrapServers = "3", SaslMechanism = SaslMechanism.ScramSha512, SaslUsername = "", SaslPassword = "", SecurityProtocol = SecurityProtocol.SaslSsl, SslCaLocation = Path.Combine(Directory.GetCurrentDirectory(), @"Certs/srdecodedca.crt"), EnableSslCertificateVerification = false }; var avroSerializerConfig = new AvroSerializerConfig { // optional Avro serializer properties: BufferBytes = 100 }; using (var schemaRegistry = new CachedSchemaRegistryClient(schemaRegistryConfig)) using (var producer = new ProducerBuilder<string, MLFixtureEmp>(producerConfig) .SetValueSerializer(new AvroSerializer<MLFixtureEmp>(schemaRegistry, avroSerializerConfig)) .Build()) { MLFixtureEmp mLFixtureEmp = new MLFixtureEmp() { deliveryDate = DateTime.Now, ISOCurrencyCode = "INR", chartererName = "ss", estimateRedeliveryDate = DateTime.Today, fixtureId = 4, imoNumber = "123", managingOwnerName = "dd", maximumDate = DateTime.Now, minimumDate = DateTime.Now, purchaseObligation = "stri", rate = 123, redeliveryPortName = "dd", RedeliveryRanges = "dsd", vesselOwnershipType = "dsds" }; producer.ProduceAsync("mytopic", new Message<string, MLFixtureEmp> { Key = "test", Value = mLFixtureEmp }).ContinueWith( task => { if (!task.IsFaulted) { Console.WriteLine($"produced to: {task.Result.TopicPartitionOffset}"); return; } Console.WriteLine($"error producing message: {task.Exception.InnerException}"); }); }
错误信息
error producing message: Confluent.Kafka.ProduceException`2[System.String,MaerskLine.CHAMPS.Dto.MLFixtureEmp]: Local: Value serialization error ---> System.InvalidOperationException: AvroSerializer only accepts type parameters of int, bool, double, string, float, long, byte[], instances of ISpecificRecord and subclasses of SpecificFixed. at Confluent.SchemaRegistry.Serdes.SpecificSerializerImpl`1..ctor(ISchemaRegistryClient schemaRegistryClient, Boolean autoRegisterSchema, Int32 initialBufferSize) at Confluent.SchemaRegistry.Serdes.AvroSerializer`1.SerializeAsync(T value, SerializationContext context) at Confluent.Kafka.Producer`2.ProduceAsync(TopicPartition topicPartition, Message`2 message, CancellationToken cancellationToken) --- End of inner exception stack trace --- at Confluent.Kafka.Producer`2.ProduceAsync(TopicPartition topicPartition, Message`2 message, CancellationToken cancellationToken)
问题原因
错误提示明确说明:AvroSerializer仅支持基础数据类型(int、bool、double等),以及实现了ISpecificRecord接口的类或SpecificFixed的子类。当前的MLFixtureEmp是自定义POCO类,未实现ISpecificRecord接口,因此无法被AvroSerializer序列化。
解决方案
方法1:使用Avro工具生成符合要求的实体类(推荐)
通过官方工具生成与Avro Schema绑定的类,确保序列化逻辑正确:
- 编写对应
MLFixtureEmp的Avro Schema文件(.avsc),注意字段类型与.NET类型的映射(例如DateTime需转为Avro的long时间戳或string格式)。 - 使用Confluent提供的
avrogen工具生成实现ISpecificRecord的实体类:avrogen -s your-schema-file.avsc ./output-directory - 在代码中替换自定义的
MLFixtureEmp为生成的类。
方法2:手动实现ISpecificRecord接口
若需保留现有类结构,需手动实现ISpecificRecord的所有成员,严格匹配Avro Schema的字段顺序和类型:
public class MLFixtureEmp : ISpecificRecord { // 原有字段定义 public DateTime deliveryDate { get; set; } public string ISOCurrencyCode { get; set; } // ... 其他字段 // 实现ISpecificRecord接口 public Schema Schema => Schema.Parse(@"{""type"":""record"",""name"":""MLFixtureEmp"",""fields"":[ {""name"":""deliveryDate"",""type"":""long""}, {""name"":""ISOCurrencyCode"",""type"":""string""}, // ... 依次添加其他字段的Schema定义 ]}"); public object Get(int fieldPos) { return fieldPos switch { 0 => new DateTimeOffset(deliveryDate).ToUnixTimeMilliseconds(), 1 => ISOCurrencyCode, // ... 其他字段的取值逻辑 _ => throw new ArgumentOutOfRangeException(nameof(fieldPos)) }; } public void Put(int fieldPos, object value) { switch (fieldPos) { case 0: deliveryDate = DateTimeOffset.FromUnixTimeMilliseconds((long)value).DateTime; break; case 1: ISOCurrencyCode = (string)value; break; // ... 其他字段的赋值逻辑 default: throw new ArgumentOutOfRangeException(nameof(fieldPos)); } } }
额外注意事项
- 确保Schema Registry中已存在对应Avro Schema,或开启自动注册(AvroSerializer默认支持自动注册)。
- 处理.NET与Avro的类型差异:例如
DateTime无原生Avro类型,需在Get/Put方法中完成转换。
内容的提问来源于stack exchange,提问作者Niranjan
相关产品推荐
相关产品推荐

