尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
kafka-examples 集成Avro与Schema Registry:KafkaAvroSerializer序列化全流程实战指南
kafka-examples 集成Avro与Schema RegistryKafkaAvroSerializer序列化全流程实战指南【免费下载链接】kafka-examplesSnippets and small examples demonstrating kafka features and configs项目地址: https://gitcode.com/gh_mirrors/kaf/kafka-exampleskafka-examples 是一个演示 Kafka 特性与配置的实战代码仓库其中 Avro 示例模块完整展示了如何集成 Avro 序列化与 Schema Registry使用 KafkaAvroSerializer 将强类型 Avro 事件写入 Kafka再用 KafkaAvroDeserializer 在消费端自动还原。本文面向新手带你走通定义 Schema → 生产 → 消费的完整链路。 为什么 Kafka 需要 Avro 序列化直接把 JSON 字符串丢进 Kafka 虽然简单但在生产环境中会有三大痛点体积大JSON 每条消息都要重复携带字段名带宽和存储成本更高无契约字段名写错、类型不匹配只有消费时才爆炸演进困难生产者加了新字段消费者老代码直接报错Avro 用二进制格式 中心化 Schema 管理解决了这些问题Schema 统一注册到 Schema Registry消息体里只存一个 5 字节的 Schema ID消费者凭 ID 取回 Schema 自动反序列化。kafka-examples仓库用一对点击流生成器 会话化消费者把这套机制讲得明明白白。 环境准备一键搭建运行环境运行 Avro 示例需要三个组件全部使用默认配置启动即可ZookeeperKafka BrokerSchema RegistryConfluent 提供默认监听http://localhost:8081克隆仓库并构建生产者git clone https://gitcode.com/gh_mirrors/kaf/kafka-examples cd kafka-examples/AvroProducerExample mvn clean package然后创建主题clicksbin/kafka-topics.sh --create --zookeeper localhost:2181 \ --replication-factor 1 --partitions 1 --topic clicks 依赖版本参考 AvroProducerExample/pom.xmlKafka 2.4.0、Confluent 5.4.0、Avro 1.7.7其中kafka-avro-serializer是核心依赖。 定义 Avro Schema一切从 LogLine.avsc 开始项目的数据模型是一个网页点击日志LogLineSchema 定义在AvroProducerExample/src/main/resources/avro/LogLine.avsc{ namespace: JavaSessionize.avro, type: record, name: LogLine, fields: [ {name: ip, type: string}, {name: timestamp, type: long}, {name: url, type: string}, {name: referrer, type: string}, {name: useragent, type: string}, {name: sessionid, type: [null,int], default: null} ] }注意namespacename组合成了 Schema Registry 中的全限定主题名JavaSessionize.avro.LogLine这是 Schema 注册与查找的关键标识。构建时avro-maven-plugin会在generate-sources阶段基于该 Schema自动生成 Java 类Specific 类代码里直接new LogLine()填充字段即可无需手写序列化逻辑。 第一步生产者——用 KafkaAvroSerializer 发送 Avro 事件核心代码在AvroProducerExample/src/main/java/com/shapira/examples/producer/avroclicks/AvroClicksProducer.java配置 Producer 只需关注 4 个参数props.put(key.serializer, io.confluent.kafka.serializers.KafkaAvroSerializer); props.put(value.serializer, io.confluent.kafka.serializers.KafkaAvroSerializer); props.put(schema.registry.url, schemaUrl); props.put(acks, all);EventGenerator同目录EventGenerator.java负责生成模拟点击事件生产者的发送逻辑非常直观LogLine event EventGenerator.getNext(); ProducerRecordString, LogLine record new ProducerRecord(clicks, event.getIp().toString(), event); producer.send(record).get();幕后发生了什么KafkaAvroSerializer首次序列化LogLine时向 Schema Registry 注册其 Schema 并拿到 ID消息体被编码为1 字节魔数 4 字节 Schema ID Avro 二进制数据同一 IP 的事件因为 Key 相同会被路由到同一分区——天然保证了同一用户的事件有序运行生产者写入 100 条点击java -cp target/uber-ClickstreamGenerator-1.0-SNAPSHOT.jar \ com.shapira.examples.producer.avroclicks.AvroClicksProducer 100 http://localhost:8081 第二步消费者——KafkaAvroDeserializer 自动反序列化消费端示例在AvroConsumerExample/src/main/java/com/shapira/examples/consumer/avroclicks/AvroClicksSessionizer.java配置与生产端对称props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, io.confluent.kafka.serializers.KafkaAvroDeserializer); props.put(schema.registry.url, url); props.put(specific.avro.reader, true);关键配置速查表配置项作用示例值key.serializer/value.serializer生产端序列化器KafkaAvroSerializerkey.deserializer/value.deserializer消费端反序列化器KafkaAvroDeserializerschema.registry.urlSchema Registry 地址http://localhost:8081specific.avro.reader用生成的 Specific 类读取trueauto.offset.reset无消费位点时的起点earliestspecific.avro.reader true表示消费时直接还原为编译期生成的LogLine类而不是通用的GenericRecord——类型安全且无需手动取字段。该消费者读取clicks主题后做了一件典型的事会话化。它用内存表记录每个 IP 的最后活跃时间间隔超过 30 分钟就递增sessionid状态管理见同目录SessionState.java再把带会话 ID 的事件写入sessionized_clicks主题——正是消费 → 加工 → 再生产的经典管道模式。运行前先创建输出主题bin/kafka-topics.sh --create --zookeeper localhost:2181 \ --replication-factor 1 --partitions 1 --topic sessions⚠️ 新手避坑清单Schema Registry 必须先启动生产者启动时会连接它注册 Schema连不上会直接抛异常namespace别乱改它决定 Schema 在 Registry 中的全名生产与消费两端必须一致Key 的序列化器是独立的本例 Key 是 IP 字符串用StringSerializer即可无需套用 Avro手动提交位点消费端关闭了自动提交auto.commit.enablefalse处理完一批后commitSync()保证不丢消息验证结果用 Confluent 自带的 Avro 控制台消费者查看落库数据bin/kafka-avro-console-consumer --zookeeper localhost:2181 --topic sessionized_clicks --from-beginning 总结通过 kafka-examples 的 Avro 生产/消费示例你掌握了完整的 Avro 序列化链路用LogLine.avsc定义 SchemaMaven 插件自动生成 Java 类生产端配置KafkaAvroSerializerschema.registry.url完成注册与编码消费端配置KafkaAvroDeserializerspecific.avro.reader自动还原强类型对象借助 Schema Registry 实现 Schema 的集中管理与版本演进这套生产 消费 Schema 治理的组合拳正是企业级 Kafka 数据管道的标准姿势。建议继续阅读仓库中的 AvroProducerExample/README.md 和 AvroConsumerExample/README.md动手跑通后再挑战 Kafka Streams 相关示例进阶流式计算 【免费下载链接】kafka-examplesSnippets and small examples demonstrating kafka features and configs项目地址: https://gitcode.com/gh_mirrors/kaf/kafka-examples创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED

相关推荐

Duster 风格指南全解读:Tighten 团队沉淀的 Laravel 代码规范精华清单

Duster 风格指南全解读:Tighten 团队沉淀的 Laravel 代码规范精华清单

Duster 风格指南全解读:Tighten 团队沉淀的 Laravel 代码规范精华清单 【免费下载链接】duster Automatic configuration for Laravel apps to apply Tightens standard linting & code standards. 项目地址: https://gitcode.com/gh_mirrors/du/duster …

📅 2026/10/6 15:05:35
RDAP与传统Whois怎么选?用ipwhois前必须知道的6个核心区别

RDAP与传统Whois怎么选?用ipwhois前必须知道的6个核心区别

RDAP与传统Whois怎么选?用ipwhois前必须知道的6个核心区别 【免费下载链接】ipwhois Retrieve and parse whois data for IPv4 and IPv6 addresses 项目地址: https://gitcode.com/gh_mirrors/ip/ipwhois RDAP与传统Whois怎么选?这是每个用 ipwho…

📅 2026/10/11 0:52:08
Zotero Connectors 如何一键抓取网页文献:零基础安装到自建完整指南

Zotero Connectors 如何一键抓取网页文献:零基础安装到自建完整指南

Zotero Connectors 如何一键抓取网页文献:零基础安装到自建完整指南 【免费下载链接】zotero-connectors Chrome, Firefox, Edge, and Safari extensions for Zotero 项目地址: https://gitcode.com/gh_mirrors/zo/zotero-connectors Zotero Connectors 是 Z…

📅 2026/10/9 6:27:29
MORE NEWS

更多资讯

📰

AI 优化内容生成是什么?从RAG引用机制到GEO落地的实战指南

一、AI 优化内容生成的本质,是让内容适配大模型的检索与引用逻辑 AI 优化内容生成(AI-Optimized Content Generation),指的是按照生成式引擎的检索增强生成(RAG)机制来组织内容,使大模型在回答用…

📰

YOLOv8注意力机制实战:SimAM、EMA、GAM源码修改与避坑指南

简介:这份学习记录面向正在使用YOLOv8做目标检测、希望借助注意力机制提升模型性能的开发者与研究者,系统整理了在YOLOv8中接入三种注意力模块的完整实践过程。内容涵盖无参数注意力SimAM、单通道注意力EMA以及双通道注意力GAM,分别给出源码引…

📰

yolov5生猪行为检测全流程:从数据集构建到训练部署实战

简介:面向养殖场智能化管理场景,YOLOv5生猪行为状态检测训练权重与PyQt界面工程,能够帮助算法工程师、农业信息化开发者快速搭建猪只进食、站立、躺卧、攻击等行为识别系统。包内包含1000多张基于养殖场视频监控帧的已标注图像及对应txt标签&…

📰

PyTorch实战样章拆解:训练循环、回归项目与DataLoader核心要点

简介:这份资源是《Deep Learning with PyTorch》的官方样章PDF,面向希望入门PyTorch框架、掌握深度学习项目实践的开发者与学习者,尤其适合具备一定Python基础、想通过动手示例理解模型训练全流程的读者。压缩包内仅含1个PDF文件,…

📰

PyTorch深度学习样本实战:从数据加载到模型训练全流程拆解

简介:这份资源是《Deep Learning with PyTorch》的官方样章PDF,面向希望入门PyTorch深度学习框架的开发者与学习者,尤其适合具备一定Python基础、想通过动手项目理解模型训练全流程的读者。样章内容围绕深度学习模型训练的核心环节展开&#…

📰

基于Open3D的点云凹凸缺陷识别:从预处理到聚类标注全流程

简介:这是一份基于Open3D的点云凹凸缺陷识别毕业论文资源,面向机器人工程、自动化检测及计算机视觉方向的本科生、研究生,也可供轨道交通装备制造相关工程技术人员参考。论文以复兴号轨道门异形曲面为对象,针对人工识别微细缺陷效…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬