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

自定义DefaultValueSerDes类实现Kafka Stream时报错求助

问题解决:Kafka Streams自定义SerDes报错(ValueSerdes为空/不兼容)

问题描述

已实现从Kafka Producer消费消息的Consumer类,尝试通过在DefaultValueSerDes中配置自定义ABC类来定制Kafka Stream,但运行出现错误。

相关代码与数据

Program.cs代码

// Stream configuration
var config = new StreamConfig();

config.ApplicationId = "app-testing";
config.BootstrapServers = "localhost:9092";
config.DefaultKeySerDes = new StringSerDes();
config.DefaultValueSerDes = new ABC();

StreamBuilder builder = new StreamBuilder();

IKStream<string, ABC> str = builder.Stream<string, ABC>("test-input");
str.Filter((k, v) => v.Data >= 25 && v.Data <= 50).To("test-output");

Topology t = builder.Build();

// Create a stream instance with topology and configuration
KafkaStream stream = new KafkaStream(t, config);

// Subscribe CTRL + C to quit stream application
Console.CancelKeyPress += (o, e) =>
{
    stream.Dispose();
};

// Start stream instance with cancellable token
await stream.StartAsync(); 

Kafka Producer发送的JSON数据

{"Name": "DND", "Data": 37}

自定义ABC.cs类代码

public class ABC : ISerDes
{
    public string Name { get; set; }
    public int Data { get; set; }

    public object DeserializeObject(byte[] data, SerializationContext context)
    {
        var bytesAsString = Encoding.UTF8.GetString(data);
        return JsonConvert.DeserializeObject<ABC>(bytesAsString);
    }

    public void Initialize(SerDesContext context)
    {
        context.Config.BootstrapServers = "localhost:9092";
        context.Config.ApplicationId = "app-testing";
    }

    public byte[] SerializeObject(object data, SerializationContext context)
    {
        var a = JsonConvert.SerializeObject(data);
        return Encoding.UTF8.GetBytes(a);
    }
}

报错信息(中文翻译)

fail: Streamiz.Kafka.Net.Processors.SourceProcessor[0]
stream-task[0|0]|processor[KSTREAM-SOURCE-0000000000]- 无法接收源数据,因为keySerdes和/或valueSerdes未设置!KeySerdes: StringSerDes | ValueSerdes: NULL
fail: Streamiz.Kafka.Net.Processors.Internal.TaskManager[0]
处理流任务0-0失败,错误如下:
Streamiz.Kafka.Net.Errors.StreamsException: stream-task[0|0]|processor[KSTREAM-SOURCE-0000000000]- 值序列化器与该处理器的实际值不兼容。请在StreamConfig中修改默认值序列化器,或者通过方法参数(使用DSL)提供正确的Serdes
在 Streamiz.Kafka.Net.Processors.AbstractProcessor2.Process(ConsumeResult2 record)
在 Streamiz.Kafka.Net.Crosscutting.ActionHelper.MeasureLatency(Action action)
在 Streamiz.Kafka.Net.Processors.StreamTask.Process()
在 Streamiz.Kafka.Net.Processors.Internal.TaskManager.Process(Int64 now)
fail: Streamiz.Kafka.Net.Processors.StreamThread[0]
stream-thread[app-testing-6991f443-2473-4c72-9e6f-b1e596edc74b-stream-thread-0]
处理过程中遇到以下错误:
Streamiz.Kafka.Net.Errors.StreamsException: stream-task[0|0]|processor[KSTREAM-SOURCE-0000000000]- 值序列化器与该处理器的实际值不兼容。请在StreamConfig中修改默认值序列化器,或者通过方法参数(使用DSL)提供正确的Serdes
在 Streamiz.Kafka.Net.Processors.AbstractProcessor2.Process(ConsumeResult2 record)
在 Streamiz.Kafka.Net.Crosscutting.ActionHelper.MeasureLatency(Action action)
在 Streamiz.Kafka.Net.Processors.StreamTask.Process()
在 Streamiz.Kafka.Net.Processors.Internal.TaskManager.Process(Int64 now)
在 Streamiz.Kafka.Net.Processors.StreamThread.<>c__DisplayClass57_0.b__3()
在 Streamiz.Kafka.Net.Crosscutting.ActionHelper.MeasureLatency(Action action)
在 Streamiz.Kafka.Net.Processors.StreamThread.Run()

错误原因分析

  1. 职责混淆:自定义ABC类同时承担了数据实体和SerDes序列化器的职责,不符合Streamiz.Kafka.Net的设计规范。ISerDes是序列化器接口,应单独实现,不应让数据类继承它。
  2. 错误配置修改:Initialize方法中错误地修改全局配置,SerDes的Initialize仅需做自身初始化,无需覆盖全局配置。
  3. 泛型接口未正确使用:使用了非泛型的ISerDes接口,而非针对实体类型的泛型ISerDes<T>,导致Streams无法识别序列化/反序列化的目标类型,进而判定ValueSerdes为空或不兼容。

解决方案

步骤1:拆分实体类与SerDes类

将ABC拆分为纯数据实体类和单独的序列化器类。

修改后ABC实体类

public class ABC
{
    public string Name { get; set; }
    public int Data { get; set; }
}

单独实现ABCSerDes序列化器

public class ABCSerDes : ISerDes<ABC>
{
    public ABC Deserialize(byte[] data, SerializationContext context)
    {
        if (data == null) return null;
        var bytesAsString = Encoding.UTF8.GetString(data);
        return JsonConvert.DeserializeObject<ABC>(bytesAsString);
    }

    public void Initialize(SerDesContext context)
    {
        // 仅做自身初始化,无需修改全局配置
    }

    public byte[] Serialize(ABC data, SerializationContext context)
    {
        if (data == null) return null;
        var json = JsonConvert.SerializeObject(data);
        return Encoding.UTF8.GetBytes(json);
    }
}

步骤2:修正StreamConfig配置

在Program.cs中使用单独的ABCSerDes作为DefaultValueSerDes:

// Stream configuration
var config = new StreamConfig();

config.ApplicationId = "app-testing";
config.BootstrapServers = "localhost:9092";
config.DefaultKeySerDes = new StringSerDes();
config.DefaultValueSerDes = new ABCSerDes(); // 替换为单独的序列化器

StreamBuilder builder = new StreamBuilder();

IKStream<string, ABC> str = builder.Stream<string, ABC>("test-input");
str.Filter((k, v) => v.Data >= 25 && v.Data <= 50).To("test-output");

Topology t = builder.Build();

// Create a stream instance with topology and configuration
KafkaStream stream = new KafkaStream(t, config);

// Subscribe CTRL + C to quit stream application
Console.CancelKeyPress += (o, e) =>
{
    stream.Dispose();
};

// Start stream instance with cancellable token
await stream.StartAsync(); 

额外说明

  • Streamiz.Kafka.Net中推荐使用泛型ISerDes<T>接口,针对具体实体类型实现序列化逻辑,确保Streams能正确识别类型匹配关系。
  • 全局配置统一在StreamConfig中设置,SerDes仅负责自身初始化,不要在其方法中修改全局配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 05:05:25