Data Engineering Zoomcamp 流式处理补充示例指南:Python Kafka、PyFlink 与 ksqlDB 实战

发布时间:2026/9/12 4:44:44
Data Engineering Zoomcamp 流式处理补充示例指南:Python Kafka、PyFlink 与 ksqlDB 实战 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 workshopIrem 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-python1.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.pyINPUT_DATA_PATH ../resources/rides.csv BOOTSTRAP_SERVERS [localhost:9092] KAFKA_TOPIC rides_json数据模型ride.pyRide类将 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__, defaultstr).encode(utf-8) }read_records()用csv.reader跳过表头后逐行构造Ridepublish_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_hooklambda 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.pyINPUT_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_deliveryself.delivery_report回调打印分区与偏移量publish()末尾调用producer.flush()确保消息真正落盘。消费者avro_example/consumer.py使用对称的AvroDeserializerdict_to_ride_record_key/dict_to_ride_record转换函数还原对象消费组 id 为datatalkclubs.taxirides.avro.consumer.2auto.offset.resetearliest。从源码可推断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 Streamingstreams-example/内含三种实现faust/基于 FaustPython 版 Kafka Streams覆盖 windowing滚动/会话窗口、branching分支、counting计数三类典型场景配套producer_taxi_json.py向datatalkclub.yellow_taxi_ride.jsontopic 持续灌入出租车 JSON 数据。pyspark/Spark Structured Streaming 消费 Kafka提供streaming.py与 Jupyter notebookstreaming-notebook.ipynb可用spark-submit.sh提交运行。redpanda/与 PySpark 示例相同仅 broker 换为 Redpanda。以 Faust 窗口聚合为例windowing.pyapp faust.App(datatalksclub.stream.v2, brokerkafka://localhost:9092) topic app.topic(datatalkclub.yellow_taxi_ride.json, value_typeTaxiRide) vendor_rides app.Table(vendor_rides_windowed, defaultint).tumbling( timedelta(minutes1), expirestimedelta(hours1), ) 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(minutes1))定义1 分钟滚动窗口、以expirestimedelta(hours1)设置窗口数据过期时间、以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 --namehadoop-distributed-file-system # 3. 启动服务分别在 kafka/ 与 spark/ 目录内执行 docker compose up -d # 4. 停止服务 docker compose downdocker/kafka/提供 Kafka 单节点 composedocker/spark/提供分层的 Spark 镜像构建cluster-base.Dockerfile、spark-base.Dockerfile、spark-master.Dockerfile、spark-worker.Dockerfile、jupyterlab.Dockerfile。注意其中docker-compose.yml与build.sh的脚本调用在源码目录中保持一致实际运行时需按各目录内的定义执行。三、pyflink/Makefile 驱动的 PyFlink 流处理 workshop3.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 查询容器化 PostgreSQLpostgres-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/ 可用后再进行下一步——首次构建镜像约需 530 分钟之后重建仅需数秒。判定集群就绪的标志是 jobmanager 日志出现taskmanager Successful registration at resource manager akka.tcp://flinkjobmanager: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_TOPICrides, VALUE_FORMATJSON);全量查询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 SELECTCTAS把窗口聚合结果物化为可查询表每条持续查询都必须以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语法想对比新旧两套 workshoppyflink/Flink 1.x Make PostgreSQL sink是旧版主课程 07-streaming 的 code/ 是基于 Flink 2.2 uv 分步 README 的重写版建议以主课程为主线、extras/为补充参考。所有示例的启动前提都是本地存在可访问的 Kafka/Redpanda brokerpython/子模块可通过docker/下的 compose 一键拉起 Kafka 与 Sparkpyflink/与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),仅供参考