自定义DefaultValueSerDes类实现Kafka Stream时报错求助
问题描述
已实现从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()
错误原因分析
- 职责混淆:自定义
ABC类同时承担了数据实体和SerDes序列化器的职责,不符合Streamiz.Kafka.Net的设计规范。ISerDes是序列化器接口,应单独实现,不应让数据类继承它。 - 错误配置修改:
Initialize方法中错误地修改全局配置,SerDes的Initialize仅需做自身初始化,无需覆盖全局配置。 - 泛型接口未正确使用:使用了非泛型的
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

