如何通过Spring-Kafka结合Confluent schema registry向Kafka发送带JSON Schema的记录
使用Spring-Kafka发送携带JSON Schema的记录实现方案
以下实现基于Spring-Kafka 2.8+、Confluent Schema Registry 7.0+版本,默认你已完成Schema Registry服务的部署。
1. 引入核心依赖
在Maven的pom.xml中添加如下依赖:
<!-- Spring-Kafka核心依赖 --> <dependency> <groupId>org.springframework.kafka</groupId> <artifactId>spring-kafka</artifactId> <version>2.9.2</version> </dependency> <!-- Confluent JSON序列化器 --> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-json-serializer</artifactId> <version>7.3.0</version> </dependency>
2. 配置生产者参数
在application.yml中添加Kafka生产者相关配置:
spring: kafka: bootstrap-servers: 你的Kafka集群地址:9092 producer: # key序列化器,根据实际业务调整 key-serializer: org.apache.kafka.common.serialization.StringSerializer # value使用Confluent提供的JSON序列化器,自动处理Schema关联 value-serializer: io.confluent.kafka.serializers.json.KafkaJsonSchemaSerializer properties: # Schema Registry服务地址 schema.registry.url: 你的Schema Registry地址:8081 # 开启自动注册Schema,生产环境如果管控严格可以关闭后手动上传Schema auto.register.schemas: true # 自动推导Java类对应的Schema类型 json.add.type.info: true
3. 定义业务数据类
编写你需要发送的业务数据对应的Java类,序列化器会自动基于该类生成JSON Schema:
// 示例业务类 public class UserMessage { private Long userId; private String userName; private Integer age; // 必需保留无参构造方法 public UserMessage() {} // 全参构造、getter、setter省略 }
4. 发送消息实现
直接注入KafkaTemplate调用发送方法即可,序列化器会自动完成Schema注册、消息头插入Schema ID的操作:
@Service public class KafkaMessageService { @Autowired private KafkaTemplate<String, UserMessage> kafkaTemplate; public void sendUserMessage(UserMessage message) { kafkaTemplate.send("目标topic名称", message.getUserId().toString(), message); } }
常见注意事项
- 生产环境如果需要管控Schema版本,可以将
auto.register.schemas设为false,提前手动将对应Schema上传到Schema Registry,同时添加配置use.latest.version: true即可使用最新版本的Schema - 如果需要自定义Schema规则,可以在Java类的字段上使用
@Schema注解指定字段说明、默认值、是否必填等属性 - Schema Registry默认的兼容性策略为BACKWARD,修改Schema时需要符合兼容性要求,否则新版本Schema会注册失败
内容的提问来源于stack exchange,提问作者rloeffel
相关产品推荐
相关产品推荐

