Data Engineering Zoomcamp 流式处理补充示例指南:Python Kafka、PyFlink 与 ksqlDB 实战
【免费下载链接】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 课程仓库cohorts/2027/07-streaming/extras/目录,系统讲解课程主线之外的一套“补充流式处理示例”(Supplementary streaming examples)。这套示例来自往期课程,虽不属于 07-streaming 主 workshop,却完整覆盖了流式处理入门的三条技术路线:基于kafka-python与confluent-kafka的 Kafka 生产者/消费者、基于 Faust 与 Spark Structured Streaming 的流处理、基于 PyFlink 与 ksqlDB 的 SQL 化流处理。读完本文,你将掌握这些示例的目录结构、依赖环境、运行方式,以及每份源码中可复用的序列化、窗口聚合与连接器配置模式。
一、示例集整体定位与目录结构
cohorts/2027/07-streaming/extras/README.md明确指出,这些示例不属于主 workshop 的一部分,而是“来自往期课程年份的额外流处理示例,可作参考资料使用”(Additional stream processing examples from previous course years)。其价值在于补充主课程未展开的多种技术栈实践,仓库当前共有三个子模块:
cohorts/2027/07-streaming/extras/ ├── README.md # 本索引文档 ├── python/ # 基于 Python 各库的 Kafka 示例(Irem Erturk 编写) ├── pyflink/ # PyFlink workshop(Irem Erturk 编写,Apache Flink 1.x) └── ksqldb/ └── commands.md # ksqlDB 查询示例需要特别说明的是模块演进关系:pyflink/目录原为 2025 年 Zach Wilson 联动的流处理内容,其后由 Alexey 基于Flink 2.2、uv 与分步 README重写为当前cohorts/2027/07-streaming/code/主 workshop(见 pyflink/README.md)。因此阅读extras/时建议同时参考主线 07-streaming 目录 以对照新旧两套实现。
二、python/:四大 Python 库的 Kafka 流处理示例
2.1 模块组成与依赖
python/README.md 将模块拆分为三大部分:
- Docker 模块:运行 Kafka 与 Spark 的 Dockerfile 与 docker-compose 定义,是所有下游示例运行的前置条件;
- Kafka 生产者/消费者示例:
json_example/(kafka-python库)与avro_example/(confluent-kafka库); - 流处理示例:
streams-example/下的 Faust、PySpark、Redpanda 三种变体。
依赖声明位于 requirements.txt,关键版本信息如下:
kafka-python==1.4.6 confluent_kafka requests avro faust fastavro安装方式为常规 pip 安装:
pip install -r cohorts/2027/07-streaming/extras/python/requirements.txt所有示例共享同一份样本数据resources/rides.csv(纽约出租车行程数据)与 Avro 模式定义,两者共同构成 python/resources/ 目录。
2.2 JSON 生产者/消费者(kafka-python)
该示例展示kafka-python最基础的用法:将 CSV 逐行读入Ride对象,以 JSON 序列化写入 Kafka,再从 Kafka 消费还原为对象。
配置集中管理(settings.py):
INPUT_DATA_PATH = '../resources/rides.csv' BOOTSTRAP_SERVERS = ['localhost:9092'] KAFKA_TOPIC = 'rides_json'数据模型(ride.py):Ride类将 CSV 18 个字段逐一映射为属性,其中时间字段用datetime.strptime解析、金额字段用Decimal保留精度、位置字段转int;同时提供from_dict类方法用于消费端反序列化。
生产者(producer.py)核心逻辑:
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') }read_records()用csv.reader跳过表头后逐行构造Ride;publish_rides()以ride.pu_location_id作为消息 key,调用producer.send()并打印record.get().offset,同时捕获KafkaTimeoutError保证长时间运行的健壮性。
消费者(consumer.py)核心配置:
config = { 'bootstrap_servers': BOOTSTRAP_SERVERS, 'auto_offset_reset': 'earliest', 'enable_auto_commit': True, 'key_deserializer': lambda key: int(key.decode('utf-8')), 'value_deserializer': lambda x: loads(x.decode('utf-8'), object_hook=lambda d: Ride.from_dict(d)), 'group_id': 'consumer.group.id.json-example.1', }consume_from_kafka()采用subscribe()+ 循环poll(1.0)的模式,将轮询超时限制为 1 秒以便响应KeyboardInterrupt,这是 Python Kafka 消费者标准写法。
运行方式(在对应示例目录下):
# 先启动生产者 python3 producer.py # 再启动消费者 python3 consumer.py2.3 Avro 生产者/消费者(confluent-kafka + Schema Registry)
与 JSON 示例不同,Avro 示例引入Schema Registry做模式管理与版本演进。配置(avro_example/settings.py):
INPUT_DATA_PATH = '../resources/rides.csv' RIDE_KEY_SCHEMA_PATH = '../resources/schemas/taxi_ride_key.avsc' RIDE_VALUE_SCHEMA_PATH = '../resources/schemas/taxi_ride_value.avsc' SCHEMA_REGISTRY_URL = 'http://localhost:8081' BOOTSTRAP_SERVERS = 'localhost:9092' KAFKA_TOPIC = 'rides_avro'生产者(avro_example/producer.py)的关键点是双 AvroSerializer 初始化:
schema_registry_client = SchemaRegistryClient({'url': props['schema_registry.url']}) self.key_serializer = AvroSerializer(schema_registry_client, key_schema_str, ride_record_key_to_dict) self.value_serializer = AvroSerializer(schema_registry_client, value_schema_str, ride_record_to_dict)发送时通过SerializationContext(topic, MessageField.KEY/VALUE)显式声明 key/value 字段,并通过on_delivery=self.delivery_report回调打印分区与偏移量;publish()末尾调用producer.flush()确保消息真正落盘。
消费者(avro_example/consumer.py)使用对称的AvroDeserializer+dict_to_ride_record_key/dict_to_ride_record转换函数还原对象,消费组 id 为datatalkclubs.taxirides.avro.consumer.2,auto.offset.reset=earliest。
从源码可推断,Avro 相对 JSON 的核心收益是:模式由.avsc文件与 Schema Registry 统一管理,key/value 各自独立定义模式(taxi_ride_key.avsc与taxi_ride_value.avsc),生产消费两端不再依赖硬编码字段顺序,天然支持模式演进而无需重写消费者。
2.4 Redpanda 变体:零配置兼容 Kafka 协议
redpanda_example/与 JSON 示例逻辑一致,仅将 broker 换成 Redpanda。该目录自带 docker-compose.yaml 与独立 README.md,意味着它可以在不依赖外部 Kafka 集群的前提下本地独立运行——Redpanda 完全兼容 Kafka 协议,示例代码无需任何修改即可切换 broker。
2.5 流处理层:Faust 与 Spark Structured Streaming
streams-example/内含三种实现:
- faust/:基于 Faust(Python 版 Kafka Streams),覆盖 windowing(滚动/会话窗口)、branching(分支)、counting(计数)三类典型场景,配套
producer_taxi_json.py向datatalkclub.yellow_taxi_ride.jsontopic 持续灌入出租车 JSON 数据。 - pyspark/:Spark Structured Streaming 消费 Kafka,提供
streaming.py与 Jupyter notebook(streaming-notebook.ipynb),可用spark-submit.sh提交运行。 - redpanda/:与 PySpark 示例相同,仅 broker 换为 Redpanda。
以 Faust 窗口聚合为例(windowing.py):
app = faust.App('datatalksclub.stream.v2', broker='kafka://localhost:9092') topic = app.topic('datatalkclub.yellow_taxi_ride.json', value_type=TaxiRide) vendor_rides = app.Table('vendor_rides_windowed', default=int).tumbling( timedelta(minutes=1), expires=timedelta(hours=1), ) @app.agent(topic) async def process(stream): async for event in stream.group_by(TaxiRide.vendorId): vendor_rides[event.vendorId] += 1这段代码演示了三个核心概念:以app.Table(...).tumbling(timedelta(minutes=1))定义1 分钟滚动窗口、以expires=timedelta(hours=1)设置窗口数据过期时间、以stream.group_by(TaxiRide.vendorId)实现按供应商 ID 的流式分组聚合——这正是 Kafka Streams 语义在 Python 中的等价实现。
2.6 Docker 环境:本地拉起 Kafka 与 Spark 集群
python/docker/README.md 给出前置环境的完整步骤:
# 1. 构建 Spark 镜像(含 JupyterLab 界面) ./build.sh # 2. 创建网络与卷 docker network create kafka-spark-network docker volume create --name=hadoop-distributed-file-system # 3. 启动服务(分别在 kafka/ 与 spark/ 目录内执行) docker compose up -d # 4. 停止服务 docker compose downdocker/kafka/提供 Kafka 单节点 compose;docker/spark/提供分层的 Spark 镜像构建(cluster-base.Dockerfile、spark-base.Dockerfile、spark-master.Dockerfile、spark-worker.Dockerfile、jupyterlab.Dockerfile)。注意其中docker-compose.yml与build.sh的脚本调用在源码目录中保持一致,实际运行时需按各目录内的定义执行。
三、pyflink/:Makefile 驱动的 PyFlink 流处理 workshop
3.1 环境依赖与安装
pyflink/README.md 声明运行前置条件为Docker(必选)、Docker Compose(必选)、Make(推荐)。Make 缺失时可手工执行Makefile中的等价命令,或按平台安装:
# Ubuntu/Debian sudo apt-get update && sudo apt-get install build-essential # CentOS/Fedora sudo dnf install make # macOS xcode-select --install # Windows(需 Chocolatey) choco install make3.2 Makefile 全量命令速查
make help输出的可用目标如下:
| 目标 | 作用 |
|---|---|
db-init | 构建并启动 PostgreSQL 数据库服务 |
build | 构建含 PyFlink 与各连接器的 Flink 基础镜像 |
up | 构建基础镜像并启动 Flink 集群 |
down | 关闭 Flink 集群 |
job | 提交 Flink 作业 |
stop/start | 停止 / 启动全部 compose 服务 |
clean | 停止并移除容器及<none>悬空镜像 |
psql | 以 psql CLI 查询容器化 PostgreSQL |
postgres-die-mac/postgres-die-pc | 清除本地挂载的 postgres 数据目录 |
3.3 三步跑通流水线
第一步:make up构建镜像并启动服务
make up # 等价命令:docker compose up --build --remove-orphans -d该步骤会一并启动 PostgreSQL 与 Flink 集群,并自动创建 sink 表processed_events。必须等到 Flink UI 在 http://localhost:8081/ 可用后再进行下一步——首次构建镜像约需 5~30 分钟,之后重建仅需数秒。判定集群就绪的标志是 jobmanager 日志出现:
taskmanager Successful registration at resource manager akka.tcp://flink@jobmanager:6123/user/rpc/resourcemanager_* under registration id <id_number>第二步:make job提交 PyFlink 作业
make job # 等价命令:docker-compose exec jobmanager ./bin/flink run -py /opt/job/start_job.py -d约一分钟后出现Job has been submitted with JobID <job_id_number>提示,即可到 Flink UI 的 Running Jobs 页面观察作业运行状态。
第三步:make stop/make down/make clean收尾清理
make stop # 停止运行中的 compose 服务 make down # 停止并移除 compose 服务 make clean # 移除容器与悬空镜像需要注意:PostgreSQL 容器的/var/lib/postgresql/data挂载到宿主机./postgres-data目录,数据在容器重启或移除后依然持久保留。
3.4 作业源码剖析:Kafka 到 PostgreSQL 的端到端管道
src/job/start_job.py是流水线的核心,通过 PyFlink Table API 以纯 SQL DDL 完成两端定义(相对路径见 start_job.py)。
Kafka 源表(此处实际指向 Redpanda,说明两者协议互通):
CREATE TABLE events ( test_data INTEGER, event_timestamp BIGINT, event_watermark AS TO_TIMESTAMP_LTZ(event_timestamp, 3), WATERMARK for event_watermark as event_watermark - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'properties.bootstrap.servers' = 'redpanda-1:29092', 'topic' = 'test-topic', 'scan.startup.mode' = 'latest-offset', 'properties.auto.offset.reset' = 'latest', 'format' = 'json' );注意这里用计算列 + 水位线声明了5 秒乱序容忍窗口(WATERMARK ... - INTERVAL '5' SECOND),这是 Flink 处理迟到事件的基石。
PostgreSQL Sink 表:
CREATE TABLE processed_events ( test_data INTEGER, event_timestamp TIMESTAMP ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:postgresql://postgres:5432/postgres', 'table-name' = 'processed_events', 'username' = 'postgres', 'password' = 'postgres', 'driver' = 'org.postgresql.Driver' );主流程在log_processing()中:以StreamExecutionEnvironment.get_execution_environment()构建环境并启用enable_checkpointing(10 * 1000)(10 秒 checkpoint),随后用一条INSERT INTO ... SELECT ... TO_TIMESTAMP_LTZ(...)将 Kafka 消息写入 PostgreSQL。src/producers/下的load_taxi_data.py与producer.py则负责向 topic 持续产生数据。
四、ksqldb/:SQL 化 Kafka 流处理
ksqldb/commands.md 提供了 ksqlDB 的查询速查手册,与 07-streaming 理论部分的 Kafka Streams 视频 配套使用。核心示例按复杂度递进:
创建流(对应 Kafka topic 的 Schema-on-Read 视图):
CREATE STREAM ride_streams ( VendorId varchar, trip_distance double, payment_type varchar ) WITH (KAFKA_TOPIC='rides', VALUE_FORMAT='JSON');全量查询:
select * from RIDE_STREAMS EMIT CHANGES;按供应商分组计数:
SELECT VENDORID, count(*) FROM RIDE_STREAMS GROUP BY VENDORID EMIT CHANGES;带过滤的分组计数:
SELECT payment_type, count(*) FROM RIDE_STREAMS WHERE payment_type IN ('1', '2') GROUP BY payment_type EMIT CHANGES;会话窗口聚合(60 秒会话,输出物化为表):
CREATE TABLE payment_type_sessions AS SELECT payment_type, count(*) FROM RIDE_STREAMS WINDOW SESSION (60 SECONDS) GROUP BY payment_type EMIT CHANGES;可以看到,ksqlDB 将 Kafka Streams 的流/表双重抽象映射为STREAM/TABLE两类对象:CREATE STREAM定义基于 topic 的实时流,CREATE TABLE ... AS SELECT(CTAS)把窗口聚合结果物化为可查询表;每条持续查询都必须以EMIT CHANGES结尾以输出增量结果。
五、如何选择与继续深入
结合主课程结构,可给出如下选型参考:
- 仅需生产/消费消息,且不在意模式管理:用
json_example/(kafka-python),代码最少、依赖最简单; - 需要强模式约束与演进能力:用
avro_example/(confluent-kafka + Schema Registry),务必保持settings.py中的SCHEMA_REGISTRY_URL指向可用的 Schema Registry 实例; - 想在 Python 生态内做有状态流处理(窗口、聚合):参考
streams-example/faust/,其App/Table/agent抽象与 Kafka Streams 一一对应; - 想要 SQL 化声明式流处理:参考
pyflink/的 Table API DDL 写法,或ksqldb/commands.md的CREATE STREAM/CTAS语法; - 想对比新旧两套 workshop:
pyflink/(Flink 1.x + Make + PostgreSQL sink)是旧版,主课程 07-streaming 的 code/ 是基于 Flink 2.2 + uv + 分步 README 的重写版,建议以主课程为主线、extras/为补充参考。
所有示例的启动前提都是本地存在可访问的 Kafka/Redpanda broker;python/子模块可通过docker/下的 compose 一键拉起 Kafka 与 Spark,pyflink/与redpanda_example/则各自内置 compose 文件,可独立运行。将上述三个子模块与主课程 07-streaming 对照研读,即可完整覆盖从消息队列到流式 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
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考