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

C# Confluent.Kafka消费者Protobuf消息自定义值反序列化器实现求助

实现方案

前置准备

  • 安装依赖NuGet包:Google.Protobuf、Confluent.Kafka
  • 用Protobuf编译器将你提供的.proto文件生成对应的C# FinalValue类,不要使用自定义的ProtobufMessage类,Protobuf生成的类自带序列化/反序列化适配逻辑:可以直接在项目中添加proto文件,将文件属性的「生成操作」设置为Protobuf 编译器,「生成输出」设置为客户端,编译项目后会自动生成对应的C#类。

第一步:实现通用Protobuf反序列化器

创建实现IDeserializer<T>接口的泛型反序列化器:

using Confluent.Kafka;
using Google.Protobuf;

public class ProtobufDeserializer<T> : IDeserializer<T> where T : IMessage<T>, new()
{
    // 缓存Parser实例,避免重复创建提升性能
    private static readonly MessageParser<T> _parser = new MessageParser<T>(() => new T());

    public T Deserialize(ReadOnlySpan<byte> data, bool isNull, SerializationContext context)
    {
        // 处理空消息场景
        if (isNull || data.IsEmpty)
        {
            return default;
        }
        return _parser.ParseFrom(data);
    }
}

第二步:替换消费者中的反序列化器配置

修改原来的ConsumerBuilder配置,替换你标记的.SetValueDeserializer行,同时修改泛型参数为生成的FinalValue类型:

using (var consumer = new ConsumerBuilder<string, FinalValue>(config)
    .SetKeyDeserializer(Deserializers.Utf8)
    // 替换原来的Utf8反序列化器为自定义Protobuf反序列化器
    .SetValueDeserializer(new ProtobufDeserializer<FinalValue>())
    .Build())

可选:映射到自定义的ProtobufMessage类

如果你需要保留原有自定义的ProtobufMessage模型,可以在消费到消息后做字段映射:

var msg = consumer.Consume(cts.Token);
var protoMsg = new ProtobufMessage
{
    timestamp = msg.Value.Timestamp.ToDateTime(),
    inputId = (int)msg.Value.InputId,
    value = msg.Value.Value?.Value,
    sourceId = (int)msg.Value.SourceId,
    inputGuid = msg.Value.InputGuid
};
Console.WriteLine($"Received: '{protoMsg.inputGuid}'");

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 20:27:02