Java中将Avro与Schema Registry结合进Kafka的核心是消息仅存schema ID而不由消息携带完整schema,通过注册表统一管理结构定义,实现空间节省、字段灵活变更、兼容性校验及多版本追溯。

在 Java 中把 Avro 和 Schema Registry 结合进 Kafka,核心是让消息不带完整 schema,只存一个 ID,靠注册表统一管理结构定义。这样既节省空间,又支持字段增删改、兼容性校验和多版本追溯。
准备 Schema Registry 服务
先确保 Confluent Schema Registry 已启动并可访问。常见做法是下载 Confluent Platform(如 7.7.x),修改 schema-registry.properties 中的监听地址和 Kafka 连接配置:
-
listeners 设为
http://0.0.0.0:8081(或内网 IP) -
kafkastore.bootstrap.servers 指向你的 Kafka broker(例如
PLAINTEXT://localhost:9092) -
kafkastore.topic 默认为
_schemas,不建议改
启动后用 curl -X GET http://localhost:8081/subjects 测试是否返回空数组,确认服务就绪。
定义 Avro Schema 并注册
写一个 user.avsc 文件(JSON 格式),比如:
立即学习“Java免费学习笔记(深入)”;
{"type":"record","name":"User","namespace":"com.example","fields":[{"name":"id","type":"long"},{"name":"name","type":"string"},{"name":"email","type":["null","string"],"default":null}]}
有两种注册方式:
在 Java 中初始化和管理阿里云 SDK客户端。包括单例模式、线程安全、endpoint 与 region 配置、VPC 终端节点、同步与异步等。
- 启动时自动注册:Producer 使用
AvroKafkaSerializer并开启auto.register.schemas=true(默认),首次发消息会自动上传 schema - 手动注册:用 curl 提前推送:
curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" --data '{"schema": "{\"type\":\"record\",...}"}' http://localhost:8081/subjects/user-value/versions
配置 Producer 和 Consumer 的 Avro Serde
依赖 Maven 引入 Confluent 官方库:
<dependency> <groupId>io.confluent</groupId> <artifactId>kafka-avro-serializer</artifactId> <version>7.7.1</version> </dependency>
Producer 配置示例:
- 序列化器设为
io.confluent.kafka.serializers.KafkaAvroSerializer - 加配置
schema.registry.url=http://localhost:8081 - 可选:关闭自动注册
auto.register.schemas=false(适合生产环境强管控场景)
Consumer 配置类似,反序列化器用 KafkaAvroDeserializer,并设置 specific.avro.reader=true(若使用生成的 Java 类)或保持默认(用 GenericRecord)。
编写业务代码(以 GenericRecord 为例)
不用生成 Java 类也能快速上手:
- Producer 构造
GenericRecord:用Schema.Parser().parse(...)加载 schema,再 newGenericData.Record(schema)赋值 - 发送时直接传 record 对象,Serde 自动处理注册与序列化
- Consumer 收到后仍是
GenericRecord,用get("fieldName")取值,无需编译类
如果项目已用 Maven Avro 插件生成了 User 类,就把 Serde 换成 SpecificAvroSerde<User>,类型更安全,IDE 支持更好。


















