Spring Cloud Stream自定义Debezium值反序列化器不生效问题排查
问题根因
你遇到的反序列化异常以及字段全为0的问题,由两个核心配置错误导致:
- Debezium Serde 配置键错误:Debezium Json Serde的配置项需要携带
debezium.json.前缀,你直接使用from.field作为配置键不被识别,Serde没有开启提取after字段的逻辑,直接尝试将完整的CDC消息反序列化为Wish类,因此报错存在未识别字段before,忽略未知属性后所有字段无法匹配,最终返回基本类型的默认值0。 - Spring Cloud Stream Kafka Streams Binder 未指定自定义Serde:你定义的Serde Bean没有被绑定到输入通道,框架默认使用全局Serde配置,你写的Serde逻辑根本没有生效。
修复方案
第一步:修正Serde配置键
修改WishSerde Bean的配置代码,补充配置前缀:
@Bean public Serde<Wish> WishSerde(){ Serde<Wish> wishSerde = DebeziumSerdes.payloadJson(Wish.class); Map<String, Object> configs = new HashMap<>(); configs.put("debezium.json.from.field", "after"); configs.put("debezium.json.unknown.properties.ignored", true); wishSerde.configure(configs, false); return wishSerde; }
第二步:绑定自定义Serde到输入通道
在application.yml中补充输入通道的Serde指定配置:
spring.cloud: function.definition: processItems stream: bindings: processItems-in-0: destination: source.wish processItems-out-0: destination: processed.wish kafka: streams: binder: brokers: 127.0.0.1:9092 bindings: processItems-in-0: consumer: keySerde: KeySerde # 对应你定义的KeySerde Bean名称 valueSerde: WishSerde # 对应你定义的WishSerde Bean名称
如果配置后提示找不到Serde,可以将上面两个配置值替换为对应Serde实现类的全限定类名。
修改完成后重启服务即可正常提取after字段内容映射为Wish对象。
内容的提问来源于stack exchange,提问作者Masoud Aghaei
相关产品推荐
相关产品推荐

