Data Engineering Zoomcamp 之 ksqlDB 流处理实战:基于 Kafka Topic 的流式 SQL 查询与窗口聚合
【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp
本文是 Data Engineering Zoomcamp 模块 7(流处理)的补充实战指南,以仓库中 07-streaming/extras/ksqldb/commands.md 为核心骨架,讲解如何用 ksqlDB 对 Kafka Topic 中的出租车乘车数据(rides)执行流式 SQL:从创建流、查询与过滤,到分组计数与会话窗口聚合。读完本文,你将掌握 ksqlDB 声明式流处理的基本模型(Stream / Table)、推送式查询(Push Query)的写法,以及如何用EMIT CHANGES持续消费实时数据。
一、ksqlDB 在本项目中的定位
在 Data Engineering Zoomcamp 的课程体系中,流处理模块(07-streaming/README.md)主线是 PyFlink Workshop(Redpanda + Python + Flink + PostgreSQL),而 Kafka 理论知识则以视频讲座形式收录在 07-streaming/theory/README.md 中,其中第 7.11 节专门讲解Kafka ksqlDB and Connect。
commands.md正是这段视频的配套代码手册。根据 07-streaming/extras/README.md 的说明,extras/目录存放的是往届课程的补充流处理示例,ksqlDB 部分被明确描述为:
example ksqlDB queries for creating streams, filtering, grouping, and windowed aggregations over Kafka topics. Companion to the ksqlDB and Connect video in the theory section.
也就是说,这份文档聚焦一个能力:用 ksqlDB 的声明式 SQL 对 Kafka Topic 中的流式数据进行建模与实时分析。它不依赖额外框架,直接在 ksqlDB 命令行或 REST 接口上执行,是理解 Kafka Streams 抽象模型(Stream/Table/窗口)最直观的入口。
二、数据模型:与仓库中的 rides 数据对应
在动手执行命令前,先明确commands.md中出现的字段含义。文档中的流模式定义了三个字段:
| 字段 | 类型 | 说明 |
|---|---|---|
VendorId | varchar | 出租车公司/供应商标识 |
trip_distance | double | 行程距离(英里) |
payment_type | varchar | 支付方式编码(字符串形式) |
这与仓库中 JSON 生产者示例的数据结构完全对应。查看 07-streaming/extras/python/json_example/ride.py 可以看到,每次乘车记录的原始 CSV 行中,vendor_id是字符串、trip_distance是数值、payment_type是字符串编码;producer.py 将其序列化为 JSON 后发送到 Kafka:
config = { 'bootstrap_servers': BOOTSTRAP_SERVERS, 'key_serializer': lambda key: str(key).encode(), 'value_serializer': lambda x: json.dumps(x.__dict__, default=str).encode('utf-8') }对应地,settings.py 中定义默认 Topic 名为rides_json、broker 地址为localhost:9092。而 ksqlDB 手册中创建的ride_streams流直接挂载在名为rides的 Topic 上,使用相同的 JSON 值格式——这正是 "Kafka 中已有 JSON 数据,用 ksqlDB 直接声明 schema 即可查询" 的典型场景。
值得注意的是,Avro 示例 中payment_type被建模为int,而 ksqlDB 手册中它是varchar且过滤条件写成IN ('1', '2')(字符串比较)。这说明在 JSON 字符串序列化方案下,支付类型字段以字符串编码在 Kafka 中流转,ksqlDB 的列类型声明必须与实际消息格式保持一致,否则查询会得到空结果或类型错误。
三、创建流:让 Kafka Topic 变成可查询的关系视图
ksqlDB 中,流(Stream)是对 Kafka Topic 的声明式封装:Topic 中的每条消息即流中的一行记录,消息 key/value 决定了行的主键与列值。文档给出的创建语句如下:
CREATE STREAM ride_streams ( VendorId varchar, trip_distance double, payment_type varchar ) WITH (KAFKA_TOPIC='rides', VALUE_FORMAT='JSON');逐步拆解:
CREATE STREAM:创建的是不可变的、仅追加的事件流视图。对应地,CREATE TABLE则用于对同一 Topic 按 key 做折叠(upsert)语义建模。- 列定义:
varchar对应字符串,double对应浮点数。列名采用大小写不敏感风格,查询时写成VENDORID与定义时的VendorId等价(下文查询语句即为全大写)。 WITH子句:KAFKA_TOPIC='rides'指定该流挂载的 Kafka Topic 名称;VALUE_FORMAT='JSON'声明消息 value 的序列化格式为 JSON。ksqlDB 会为未声明的 key 自动生成虚拟列ROWKEY。- 执行方式:在 ksqlDB CLI(
docker exec -it <ksqldb-container> ksql)或 HTTP 接口中执行该语句即可持久化注册流定义。
前提约束:Kafka 集群与 ksqlDB 需已就绪,且
ridesTopic 已存在(或允许 ksqlDB 自动创建)。本仓库的流处理环境通常由 07-streaming/extras/python/docker 下的 Docker Compose 编排 Kafka 与相关组件。
四、推送式查询:用 EMIT CHANGES 实时读取数据
流是有生命的:只要 Topic 持续有新消息,查询结果就持续更新。文档中的第一条查询:
select * from RIDE_STREAMS EMIT CHANGES;关键点在EMIT CHANGES:
- 这是推送式查询(Push Query)的标志——查询不会返回一个静态结果集就结束,而是持续订阅流,每当有新消息进入便输出一行;
- 不带
EMIT CHANGES的普通SELECT则是拉取式查询(Pull Query),仅对已物化的表状态返回当前快照,适用于CREATE TABLE ... AS SELECT产出的物化视图; - 对于流式
SELECT *,输出顺序与 Topic 中消息到达顺序一致,每条消息一行,字段按声明 schema 展开。
五、分组聚合:实时统计每个 Vendor 的乘车量
流的下一层威力在于聚合。文档第二条查询:
SELECT VENDORID, count(*) FROM RIDE_STREAMS GROUP BY VENDORID EMIT CHANGES;GROUP BY VENDORID将流按VENDORID分组,count(*)统计组内累积的记录数;- 因为是流式聚合,ksqlDB 在内存/状态存储中为每个
VENDORID维护一个计数,每条新消息到达都会触发一次增量更新并输出最新计数; - 聚合产生的内部分区 Topic 以分组键(此处为
VENDORID)为 key,保证同一供应商的事件路由到同一分区,从而保证计数一致; - 注意
VENDORID的大写写法与创建时VendorId等价,体现 ksqlDB 标识符的大小写不敏感性。
六、带过滤的聚合:只统计关心的支付方式
实时分析往往要先过滤再聚合。文档第三条查询:
SELECT payment_type, count(*) FROM RIDE_STREAMS WHERE payment_type IN ('1', '2') GROUP BY payment_type EMIT CHANGES;WHERE在聚合前对事件做过滤,只让payment_type为'1'(信用卡)或'2'(现金)的记录进入计数;- 过滤发生在流进入聚合算子之前,因此未被选中的支付方式根本不会出现在结果中,也不会占据状态存储;
- 这与前文所述数据模型呼应:在 JSON 字符串序列化方案下,支付类型以字符串形式比较,与 ride.py 中
payment_type = arr[9](原始 CSV 字符串)一致。
七、窗口函数:会话窗口内的滚动统计
窗口是流处理的灵魂,它把无界流切成有界的时间段做聚合。文档第四条查询演示了会话窗口(SESSION window):
CREATE TABLE payment_type_sessions AS SELECT payment_type, count(*) FROM RIDE_STREAMS WINDOW SESSION (60 SECONDS) GROUP BY payment_type EMIT CHANGES;要点解析:
WINDOW SESSION (60 SECONDS):会话窗口以不活动间隔为界——如果某个payment_type超过 60 秒没有新事件,则当前窗口关闭,下一个事件开启新窗口。窗口长度不是固定的,而是由数据活跃度动态决定,适合用户行为会话、连续乘车等场景;CREATE TABLE ... AS(CTAS):将窗口化聚合结果物化为一张持续更新的表。会话窗口有明确的开始/结束时间,窗口本身成为主键的一部分(与ROWKEY、WINDOWSTART、WINDOWEND一起构成表的键);EMIT CHANGES出现在 CTAS 语句中,表示该物化表通过推送模式对外输出更新;此后对该表的SELECT拉取查询可以拿到当前窗口快照;- 会话窗口之外,ksqlDB 还支持固定长度的TUMBLING(滚动窗口)与HOPPING(跳跃窗口),分别适用于固定周期统计与滑动时间窗分析。
窗口化聚合是 ksqlDB 相对普通 SQL 的最大差异点:它把 "状态" 内置到了 SQL 语法中,开发者无需手动管理状态存储与水位线。
八、从命令到源码:仓库中的配套证据链
这份命令手册不是孤立文档,仓库中与之配套的实现与说明包括:
- 07-streaming/extras/README.md — 明确 ksqlDB 命令手册是 theory 部分 ksqlDB/Connect 视频的配套材料,属于"补充参考"性质;
- 07-streaming/theory/README.md — 列出 Kafka Streams 系列讲座(含 ksqlDB 与 Schema Registry 主题),Java 示例位于 07-streaming/theory/java/kafka_examples;
- 07-streaming/extras/python/json_example — 提供与
ride_streamsschema 同源的 JSON 数据生产者与消费者,可作为向ridesTopic 灌入测试数据的参考实现; - 07-streaming/extras/python/avro_example — 展示字段类型建模的另一种选择(
payment_type为 int),可与 ksqlDB 的 varchar 声明对照理解类型匹配的重要性。
如需继续深入,原始commands.md中还给出了 ksqlDB 官方参考文档与 Java 客户端文档的入口,本文不再赘述外部链接;核心结论是:ksqlDB 让你用纯 SQL 完成原本需要编写 Kafka Streams 处理器才能实现的流式过滤、分组与窗口聚合,非常适合作为学习流处理抽象模型的起点,以及在生产中用最小代码量对 Topic 做实时探查。
九、小结:一条完整的 ksqlDB 学习路径
对照本仓库的流处理模块,建议按如下顺序消化这份命令手册:
- 用 json_example 的生产者向 Kafka 写入 JSON 格式的乘车数据(字段含
VendorId、trip_distance、payment_type); - 启动 ksqlDB,按本文第二节创建
ride_streams流; - 依次执行第四节到第七节的四种查询,观察
EMIT CHANGES下的持续输出; - 最后尝试把
SESSION窗口换成TUMBLING/HOPPING窗口,体会不同窗口语义对聚合结果形态的影响。
掌握这四个层次(建流 → 推送查询 → 分组聚合 → 窗口聚合),你就具备了用 ksqlDB 对任意 Kafka Topic 进行实时 SQL 分析的基本能力,也为后续理解 Kafka Streams 的状态管理与窗口机制打下基础。
【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考