
掌握Kafka Streams API:基于Java与Spring Boot 3.x搭建专业级实时Kafka流处理应用,从底层原理到企业级实战完整落地流处理开发方案。
课程核心学习内容
- 借助Streams API搭建完整高级Kafka Streams流处理程序
- 运用官方高阶DSL开发标准化Kafka Streams业务应用
- 实现Exactly Once精准语义:事务机制与幂等生产者落地
- 基于Streams API开发零售行业实时数据流处理项目
- 多条原始事件数据聚合生成统一聚合事件
- 多异构数据流合并为单一统一处理流
- 基于时间窗口分组对流数据完成聚合运算
- 依托Spring Boot开发符合企业规范的Kafka Streams服务
- 通过TopologyTestDriver搭配JUnit5完成Kafka Streams单元测试
- 使用EmbeddedKafka结合JUnit5实现Spring Kafka Streams集成测试
- 开发交互式查询接口,通过RESTful API读取聚合状态数据
- 多实例分布式部署下的Kafka Streams交互式查询(微服务架构模式)
课前学习要求
- 具备扎实的Java基础开发能力
- 拥有Kafka基础应用开发实战经验
- 熟练使用IntelliJ IDEA或其他主流开发IDE
- 本地开发环境配置Java 17
- 掌握Gradle或Maven项目构建工具基础用法
课程简介
Kafka Streams API是Kafka生态内置的原生高阶流处理开发接口,无需额外集群部署即可实现实时数据流转处理。
依托Kafka Streams API,开发者可一站式完成以下核心数据处理能力:
- 实时数据流清洗、格式转换与数据加工
- 外部数据关联,完成原始业务数据丰富化处理
- 单一流拆分分支,多路并行数据处理
- 跨主题多源数据合并、聚合统计计算
- 基于时间窗口的数据分桶聚合统计等场景
《使用Java/SpringBoot开发Kafka Streams API》实战课程,完整覆盖Streams API理论知识与工程编码实操,重点讲解结合Spring Boot快速搭建企业标准化Kafka流处理应用的全套技术方案。
本课程全程实战导向,所有技术知识点配套完整可运行代码,跟随实操即可吃透底层原理。课程收尾阶段,你将独立完成一套可上线运行的实时Kafka Streams业务应用。
课程学习完毕后,你将彻底吃透以下核心技术体系:
- 原生Streams API独立开发Kafka Streams流处理程序
- Spring Boot整合Streams API快速落地流处理服务
- 开发交互式查询逻辑,读取状态存储聚合数据并对外暴露RESTful接口
- 基于JUnit5完成Kafka Streams应用单元测试与全链路集成测试
模块一:Kafka Streams快速入门
本章节系统讲解Kafka Streams核心概念,梳理开发流处理应用过程中涉及的全部专业术语。
- Kafka Streams整体功能与应用场景概述
- Kafka Streams核心术语、底层架构与处理器组件详解
- KStreams基础API功能与使用场景总览
模块二:KStreams API入门——构建首个流处理问候应用
本章节从零搭建简易Kafka Streams示例项目,完成本地环境全流程调试运行。
- 设计并编写基础问候业务流拓扑结构
- 封装应用启停工具类,实现Kafka Streams程序灵活启动与关闭
模块三:KStream API内置算子详解
全面拆解Kafka Streams API提供的各类数据转换算子,掌握数据流加工核心手段。
- Filter、FilterNot:数据过滤算子
- Map、MapValues:单条记录字段转换算子
- FlatMap、FlatMapValues:一对多数据拆分算子
- peek:数据流日志打印、观测调试算子
- merge:多数据流合并算子
模块四:Kafka Streams序列化与反序列化机制
深入解析流处理中消息序列化、反序列化底层运行逻辑,解决消息传输格式兼容问题。
- Kafka Streams键值序列化、反序列化底层运行原理
- 全局默认序列化器、反序列化器配置方案
- 自定义序列化/反序列化类,扩展问候消息数据格式解析能力
模块五:通用可复用序列化器最佳实践
讲解通用泛型序列化、反序列化工具类开发方案,一套代码适配全部业务消息类型,降低重复开发成本。
- 通用泛型序列化、反序列化工具类完整实现
模块六:实战案例——零售订单管理Kafka Streams应用
基于真实零售业务场景,完整落地一套订单实时流处理系统,串联前文全部基础API知识点。
模块七:Kafka Streams内部运行原理——拓扑、流与任务
拆解Kafka Streams底层执行机制,理解拓扑、数据流、任务分区的调度运行逻辑。
- 拓扑、流、任务底层内部执行机制深度解析
模块八:Kafka Streams异常与错误处理体系
覆盖流处理程序全链路异常捕获方案,解决消息解析、处理、发送阶段报错导致服务中断问题。
- Kafka Streams全局故障处理流程
- 默认反序列化异常处理行为
- 自定义反序列化异常处理器开发
- 内置处理器默认异常策略与自定义异常处理器
- 消息生产阶段自定义错误捕获逻辑
模块九:KTable与全局GlobalKTable详解
解析有状态流处理两大核心组件KTable、GlobalKTable,掌握维度数据存储与实时更新方案。
- KTable核心API基础介绍
- 基于KTable构建状态存储拓扑流程
- KTable底层数据更新、存储机制拆解
- 全局维度表GlobalKTable适用场景与开发
模块十:有状态流操作:聚合、关联与时间窗口
系统学习Kafka Streams有状态运算能力,实现统计聚合、多流关联、分时窗口计算等核心业务需求。
- Kafka Streams有状态算子整体介绍
- 聚合运算底层原理,count计数算子实操
- groupBy分组算子使用规范
- reduce归约聚合算子业务落地
- aggregate自定义聚合逻辑开发
- 物化视图实现count、reduce持久化统计结果
模块十一:状态运算结果读取方案
讲解状态存储数据读取逻辑,掌握如何获取聚合统计后的业务计算结果。
模块十二:流数据重分区与重新键值化
实操演示空值算子使用场景,说明有状态运算中记录重键、重分区的必要性与实现代码。
模块十三:有状态操作——多数据流Join关联运算
覆盖Kafka Streams全部流关联算子,实现多源业务数据实时拼接整合。
模块十四:实战落地——订单管理系统多流Join关联
在零售订单项目中集成各类Join关联逻辑,还原真实业务多表实时关联场景。
- Kafka Streams各类Join关联类型基础介绍
- KStream与KTable innerJoin内连接实操
- KStream与GlobalKTable innerJoin内连接实操
- KTable与KTable innerJoin内连接实操
- KStream与KStream innerJoin内连接实操
- leftJoin左连接业务实现
- outerJoin全连接业务实现
- Join关联底层运行机制拆解
- 流关联前置CoPartitioning同分区要求与原理
模块十五:有状态操作——时间窗口计算
详解窗口化分时统计能力,实现分时销量、分时订单量等时段统计需求。
- 窗口计算与事件时间、处理时间基础概念
- 滚动窗口Tumbling Window开发实操
- supress算子控制窗口统计结果输出时机
- 滑动窗口Sliding Window使用场景与代码实现
模块十六:实战案例——订单系统窗口分时统计
基于零售订单项目新增分时统计需求,完整落地窗口聚合业务代码。
模块十七:窗口内时间乱序消息处理策略
分析携带滞后时间戳、未来时间戳消息在窗口计算中的处理逻辑与表现。
模块十八:Spring Boot整合Kafka Streams基础开发
快速上手Spring Boot生态下的Kafka流处理应用,依托自动配置简化开发流程。
- Spring Boot与Kafka Streams整合基础介绍
- 搭建Spring Kafka Streams问候示例项目
- application.yml配置文件统一管理流处理参数
- 业务拓扑编写与组装
- 本地环境完整调试测试问候应用
模块十九:Spring Boot Kafka Streams自动配置原理
拆解Spring Starter内置自动化装配逻辑,理解容器如何自动初始化Streams环境。
模块二十:Spring Kafka Streams JSON序列化方案
基于Spring Boot实现JSON格式消息的序列化与反序列化工具封装。
模块二十一:Spring Kafka Streams全链路异常处理
针对Spring体系下的流处理服务,提供多套异常捕获解决方案,保障服务稳定性。
- 方案一:通用反序列化异常拦截处理
- 方案二:自定义全局反序列化异常处理器
- 方案三:Spring专属内置异常处理组件
- 拓扑流程中未捕获全局异常统一拦截
- 消息发送生产阶段异常捕获处理
模块二十二:企业级实战——Spring Boot零售订单流项目搭建
完整搭建基于Spring Boot的零售订单实时流处理工程,整合前文全部Spring整合技术点。
模块二十三:交互式查询开发——RESTful接口读取状态存储
开发Web查询接口,对外暴露聚合统计数据,支持前端、其他微服务实时读取流计算结果。
- 接口开发第一部分:按订单类型统计订单数量GET接口
- 接口开发第二部分:完善订单类型统计查询能力
- 按订单类型+地区ID多条件统计订单数量接口
- 全量订单类型总数量查询通用接口
- 按订单类型统计营收金额查询接口
- 全局统一异常返回封装,输出标准化客户端错误信息
模块二十四:交互式查询开发——窗口状态存储REST接口
针对窗口分时聚合数据开发专属查询接口,支持分时统计数据实时读取。
- 按订单类型查询窗口时段订单数量接口
- 全部订单类型分时总订单量查询接口
- 自定义时间区间滑动窗口订单统计查询接口
- 按订单类型分时营收金额查询接口
模块二十五:单元测试——TopologyTestDriver + JUnit5
讲解无依赖本地单元测试方案,无需启动Kafka服务即可验证拓扑业务逻辑正确性。
- TopologyTestDriver基础使用方法
- 问候应用单元测试:校验输出主题消息
- 多批次消息流转场景单元测试
- 异常报错场景覆盖单元测试
- 订单计数状态存储读写单元测试
- 订单营收统计状态存储单元测试
- 多分区订单营收数据单元测试
- TopologyTestDriver工具局限性说明
模块二十六:Spring Boot项目单元测试集成TopologyTestDriver
结合Spring容器环境,为Spring Boot整合的Kafka Streams服务编写完整单元测试用例。
模块二十七:集成测试——EmbeddedKafka模拟Kafka服务
依托嵌入式Kafka组件完成全链路集成测试,模拟真实消息收发环境校验业务流程。
- 嵌入式Kafka集成测试介绍与环境搭建
- 订单数量统计全链路集成测试
- 订单营收金额集成测试用例
- 分时窗口营收统计集成测试
模块二十八:Kafka Streams消息宽限期配置
讲解窗口计算宽限期核心概念,配置延迟消息保留时长,提升分时统计数据完整性。
模块二十九:Spring Boot项目打包与生产环境部署
完整流程实现流处理应用打包为可执行Jar包,掌握线上服务启动运维方式。
课程学习收获
完成本课程全部章节学习后,你将深度吃透Kafka Streams API完整技术栈,独立开发各类场景下的实时流处理应用,满足企业生产环境落地标准。
课程适配人群
- 有多年开发经验的高级Java工程师
- 掌握Kafka基础,希望深耕Kafka Streams流处理技术的开发人员
- 需要搭建高性能实时流处理系统的Kafka开发工程师
- 想要掌握TopologyTestDriver自动化测试流应用的后端技术开发者
