Spring Cloud Stream Kafka Topic与Schema Registry Subject不匹配问题求助
问题根因
Confluent官方的Avro序列化/反序列化器默认使用TopicNameStrategy作为Subject命名策略,该策略会直接将Topic完整名称拼接-key/-value后缀作为Subject名称请求Schema Registry,因此你Topic为abc.bla时就会请求abc.bla相关的Subject,和你预期的abc不符。
解决方案
自定义Subject命名策略适配需求
你需要实现符合自己业务规则的Subject命名策略,实现SubjectNameStrategy接口,逻辑为截取Topic名称第一个.之前的部分作为Subject名:
import io.confluent.kafka.serializers.subject.SubjectNameStrategy; public class CustomTopicPrefixSubjectStrategy implements SubjectNameStrategy { @Override public String subjectName(String topic, boolean isKey, Object schema) { // 取topic第一个.之前的部分作为subject,可根据你的实际规则调整 String subjectPrefix = topic.split("\\.")[0]; // 如果你的Schema Registry侧的Subject带-key/-value后缀,可保留下方拼接逻辑 // return subjectPrefix + (isKey ? "-key" : "-value"); // 你的需求直接返回前缀即可 return subjectPrefix; } }
修改Spring配置加载自定义策略
将自定义的策略类配置到序列化/反序列化参数中,修改你的application.yaml对应部分:
spring: cloud: stream: kafka: binder: consumer-properties: # 新增value反序列化的subject策略配置,替换为你自定义类的全限定名 value.subject.name.strategy: com.xxx.config.CustomTopicPrefixSubjectStrategy producer-properties: # 新增value序列化的subject策略配置,替换为你自定义类的全限定名 value.subject.name.strategy: com.xxx.config.CustomTopicPrefixSubjectStrategy
如果需要同时调整Key对应的Subject命名规则,额外添加key.subject.name.strategy配置即可。配置生效后重启服务,请求Schema Registry的Subject就会按照你定义的规则生成,匹配预期的abc名称。
内容的提问来源于stack exchange,提问作者Clueless
相关产品推荐
相关产品推荐

