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

.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
 ---&gt; 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绑定的类,确保序列化逻辑正确:

  1. 编写对应MLFixtureEmp的Avro Schema文件(.avsc),注意字段类型与.NET类型的映射(例如DateTime需转为Avro的long时间戳或string格式)。
  2. 使用Confluent提供的avrogen工具生成实现ISpecificRecord的实体类:
    avrogen -s your-schema-file.avsc ./output-directory
    
  3. 在代码中替换自定义的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 22:37:02