
Data Engineering Zoomcamp使用 PyFlink 搭建 Kafka 到 PostgreSQL 的实时流处理管道【免费下载链接】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 仓库中 pyflink 实训模块 为核心完整讲解如何基于 Apache Flink 1.16 与 PyFlink 构建一条「Kafka 兼容消息队列Redpanda→ Flink 流式作业 → PostgreSQL」的实时数据处理管道。读完本文你将掌握容器化 Flink 集群JobManager TaskManager的搭建方式、PyFlink 作业的提交与调度方法、基于 Table API 的 Kafka 数据源与 JDBC 汇入目标的定义以及滚动窗口聚合与水位线Watermark的实际用法并能在本地一键复现整条流处理链路。一、模块概览与架构pyflink是 Data Engineering Zoomcamp 第七周「流处理」课程的扩展实训模块位于 cohorts/2027/07-streaming/extras/pyflink/。它的目标是用最小的工程代价跑通一条完整的流处理流水线目录结构如下pyflink/ ├── Dockerfile.flink # Flink Python 3.7 PyFlink 各连接器的镜像定义 ├── docker-compose.yml # Redpanda、JobManager、TaskManager、PostgreSQL 编排 ├── Makefile # 封装了构建、启停、提交作业等常用命令 ├── requirements.txt # PyFlink 与 Python 依赖 ├── README.md # 本模块的完整运行指南 ├── homework.md # 配套作业 └── src/ ├── job/ # PyFlink 流式作业源码 │ ├── start_job.py # Kafka → Postgres 透传作业 │ ├── aggregation_job.py # 滚动窗口聚合作业 │ └── taxi_job.py # 出租车数据入库作业 └── producers/ # 消息生产端 ├── producer.py # 向 test-topic 发送测试消息 └── load_taxi_data.py # 将出租车 CSV 数据灌入 green-data topic整套架构由四个核心服务组成见 docker-compose.yml服务镜像端口职责redpanda-1redpandadata/redpanda:v24.2.189092 / 29092 / 8082 / 28082Kafka 协议兼容的消息队列承担 Kafka 角色jobmanagerpyflink:1.16.0本地构建8081Flink UIFlink 作业管理与调度提交 PyFlink 作业的入口taskmanagerpyflink:1.16.06121 / 6122执行流处理算子默认 15 个任务槽、并行度 3postgrespostgres:145432流处理结果的落地目标数据流方向为生产者把消息写入 Redpanda topic → PyFlink 作业从 topic 消费 → 经流式处理透传或窗口聚合→ 写入 PostgreSQL 表。二、环境准备Docker、Docker Compose 与 Make运行本模块需要三个前置组件对应 README.md 中的 Installation 部分Docker必需——承载 Flink 集群、Redpanda 与 PostgreSQLDocker Compose必需——用于一次拉起全部服务Make推荐——Makefile封装了所有常用命令使用make可以让操作更简洁。如果你的系统没有安装 Make也可以把 Makefile 中的命令复制到终端手动执行。Make 的安装方式因操作系统而异# 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 make其中build-essential除了提供make还包含 GCC 等编译工具这也是后续在容器内从源码编译 Python 3.7 时所需的依赖。确认环境后进入模块目录原 README 中写作cd 07-streaming/pyflink对应仓库内的完整相对路径是 cohorts/2027/07-streaming/extras/pyflinkcd cohorts/2027/07-streaming/extras/pyflink三、一键拉起构建镜像并启动 Flink 集群3.1 make up 做了什么在模块根目录执行make up等价于手动执行见 Makefiledocker compose up --build --remove-orphans -d这一条命令完成三件事依据 Dockerfile.flink 构建名为pyflink:1.16.0的 Flink 基础镜像依次启动 Redpanda、JobManager、TaskManager、PostgreSQL 四个服务创建 PostgreSQL 中的 sink 表由 PyFlink 作业中的CREATE TABLE语句在作业启动时完成。3.2 镜像构建的关键点Dockerfile.flink 的构建逻辑值得拆解基础镜像flink:1.16.0-scala_2.12-java8即 Flink 1.16.0 官方镜像使用 Scala 2.12 与 Java 8Python 3.7 从源码编译Debian 11 默认 Python 为 3.9而 PyFlink 1.16 官方仅支持 Python 3.6 / 3.7 / 3.8因此镜像内通过wget下载 Python 3.7.9 源码并./configure --enable-shared编译安装随后建立python软链接PyFlink 安装通过pip3 install -r requirements.txt安装 requirements.txt 中声明的依赖apache-flink1.16.0、psycopg2-binary2.9.1、requests、kafka-python连接器 Jar 下载从 Maven 中央仓库下载并放入/opt/flink/lib/这是 SQL/Table API 能读写 Kafka 与 PostgreSQL 的关键flink-json-1.16.0.jar—— JSON 格式解析flink-sql-connector-kafka-1.16.0.jar—— KafkaRedpanda连接器flink-connector-jdbc-1.16.0.jar—— JDBC 连接器postgresql-42.2.24.jar—— PostgreSQL 驱动内存参数向flink-conf.yaml追加taskmanager.memory.jvm-metaspace.size: 512m为 JVM 元空间预留足够内存避免运行期 OOM。⚠️首次构建耗时提示第一次构建镜像需要从源码编译 Python 并下载多个依赖通常需要5 到 30 分钟只要不删除本地镜像后续重建通常只需几秒钟。3.3 确认集群就绪镜像构建完成后Docker 会自动启动 JobManager 与 TaskManager这一步通常需要一分钟左右。可以在 Docker Desktop 中查看容器日志当 TaskManager 日志中出现类似下面这行时说明集群已就绪taskmanager Successful registration at resource manager akka.tcp://flinkjobmanager:6123/user/rpc/resourcemanager_* under registration id id_number然后访问 Flink Web UIhttp://localhost:8081/看到 UI 正常渲染后再进入下一步。3.4 docker-compose.yml 中的关键配置在 docker-compose.yml 中有几处直接影响运行结果的配置Redpanda 双地址监听PLAINTEXT://redpanda-1:29092供容器内服务Flink 作业使用OUTSIDE://localhost:9092供宿主机上的生产者使用——这正是src/producers/producer.py里bootstrap_serverslocalhost:9092能连通的原因JobManager 环境变量通过POSTGRES_URL、POSTGRES_USER、POSTGRES_PASSWORD、POSTGRES_DB注入 PostgreSQL 连接信息且带默认值postgres可用.env文件覆盖FLINK_PROPERTIES中声明jobmanager.rpc.address: jobmanagerTaskManager 并行配置taskmanager.numberOfTaskSlots: 15、parallelism.default: 3即 15 个任务槽、默认并行度 3目录挂载./src/:/opt/src把作业源码挂载进容器./:/opt/flink/usrlib挂载整个模块目录注意 PostgreSQL 容器没有挂载数据卷数据持久化策略见下文「清理与数据持久化」host.docker.internal通过extra_hosts映射使容器内可以访问宿主机上的服务。四、提交并运行 PyFlink 作业4.1 提交透传作业集群就绪后提交第一个 PyFlink 作业make job等价于手动执行Makefiledocker compose exec jobmanager ./bin/flink run -py /opt/src/job/start_job.py --pyFiles /opt/src -d命令拆解docker compose exec jobmanager进入 JobManager 容器执行命令./bin/flink runFlink 命令行提交作业-py /opt/src/job/start_job.py指定要运行的 Python 作业文件挂载自仓库 src/job/start_job.py--pyFiles /opt/src把整个src目录加入 Python 依赖路径供作业内 import 使用-ddetached 模式提交后立即返回作业在后台运行。大约一分钟后终端会出现Job has been submitted with JobID job_id_number的提示。此时回到 Flink UI 的 Running Jobs 页面即可看到作业在运行。4.2 产生测试数据打开另一个终端运行生产者向 Redpanda 的test-topic写入测试消息python src/producers/producer.pyproducer.py 的核心逻辑是用kafka-python的KafkaProducer连接localhost:9092循环 990 次range(10, 1000)每条消息为{test_data: i, event_timestamp: time.time() * 1000}间隔 50ms 发送最后flush()确保全部落盘。event_timestamp使用毫秒级时间戳与作业中的TO_TIMESTAMP_LTZ(event_timestamp, 3)解析逻辑对应3 表示毫秒精度。4.3 验证数据入库start_job.py定义了两张表Kafka 数据源topictest-topicCREATE 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 );PostgreSQL 汇入目标表名processed_eventsCREATE 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 );作业通过一条INSERT INTO ... SELECT完成透传见 start_job.py把 Kafka 中的BIGINT时间戳转换为 PostgreSQL 的TIMESTAMP后写入。可以看到数据源定义里有一个水位线Watermark声明WATERMARK for event_watermark as event_watermark - INTERVAL 5 SECOND它允许事件时间最多迟到 5 秒为后续基于事件时间的窗口聚合做准备。可以用make psql进入 PostgreSQL CLI 直接查询验证SELECT * FROM processed_events ORDER BY event_timestamp DESC LIMIT 10;4.4 窗口聚合作业进阶在make job之上Makefile 还提供了第二个作业入口make aggregation_job对应的 aggregation_job.py 演示了基于事件时间的滚动窗口Tumbling Window聚合这是流处理中最常用的模式之一其关键差异点在于数据源偏移策略scan.startup.mode earliest-offset、auto.offset.reset earliest即从 topic 最早的消息开始消费以便对历史数据做完整聚合水位线event_watermark - INTERVAL 1 SECOND聚合 SQL使用TUMBLE表值函数按 1 分钟窗口分组INSERT INTO processed_events_aggregated SELECT window_start as event_hour, test_data, COUNT(*) AS num_hits FROM TABLE( TUMBLE(TABLE events, DESCRIPTOR(event_watermark), INTERVAL 1 MINUTE) ) GROUP BY window_start, test_data;汇入目标processed_events_aggregated表包含event_hour、test_data、num_hits三个字段并声明PRIMARY KEY (event_hour, test_data) NOT ENFORCEDJDBC 汇入目标支持按主键 upsert。运行该作业前需要先启动生产者并且聚合作业消费的是earliest偏移因此会把历史上已发送到test-topic的消息一并按分钟窗口统计得到类似「某分钟窗口内某个 test_data 值出现了多少次」的聚合结果。4.5 出租车真实数据作业模块还附带了一个面向真实业务数据的作业 taxi_job.py数据源Kafka topicgreen-data消费模式earliest-offset格式为 JSON包含 20 个出租车行程字段VendorID、lpep_pickup_datetime、trip_distance、total_amount等并通过TO_TIMESTAMP(lpep_pickup_datetime, yyyy-MM-dd HH:mm:ss)计算事件时间水位线为 15 秒生产端load_taxi_data.py 读取data/green_tripdata_2019-10.csv需自行下载放置可参考同仓库 06-batch 的下载脚本用csv.DictReader逐行转为字典后发送到green-datatopic汇入目标PostgreSQL 表taxi_events使用CREATE OR REPLACE TABLE保证可重复执行。五、Make 命令速查表运行make help可查看全部可用命令见 Makefile当前支持的目标如下目标作用make help显示帮助信息make db-init构建并运行 PostgreSQL 数据库服务make build构建内置 PyFlink 与连接器的 Flink 基础镜像make up构建镜像并启动整个 Flink 集群含 PostgreSQLmake down关闭并移除 Flink 集群Compose 服务make job提交透传 PyFlink 作业make aggregation_job提交滚动窗口聚合作业make stop停止 Docker Compose 中的所有服务make start启动 Docker Compose 中的所有服务make clean停止并移除容器同时清理none标签的悬空镜像make psql在命令行中查询容器化 PostgreSQL 数据库make postgres-die-mac删除本机macOS挂载的 postgres 数据目录及容器内数据make postgres-die-pc删除本机PC挂载的 postgres 数据目录及容器内数据六、停止、清理与数据持久化运行结束后按需选择清理命令make stop # 停止 Docker Compose 中运行的服务 make down # 停止并移除 Docker Compose 服务 make clean # 移除容器及悬空镜像注意本模块的 PostgreSQL 容器将/var/lib/postgresql/data容器内数据目录挂载到宿主机./postgres-data目录。因此即使容器被停止或删除已写入的数据仍会持久保留在本地不会丢失。若确实需要彻底清空数据可使用make postgres-die-mac或make postgres-die-pc按宿主机平台选择它们会同时删除本机挂载数据目录和容器内数据。七、常见问题与排查思路Flink UI 打不开确认make up已执行完成且 JobManager 容器处于运行状态首次构建需要 5–30 分钟期间 UI 不可用属正常现象。作业提交后立即失败检查 TaskManager 日志中是否出现Successful registration at resource manager作业需要等待集群注册完成才能正常调度。生产者连接失败宿主机生产者必须使用localhost:9092Redpanda 的 OUTSIDE 监听而容器内 Flink 作业使用redpanda-1:29092PLAINTEXT 监听两者不能混用。聚合作业没有数据确认生产者已先启动并向test-topic发送了消息且聚合作业使用earliest-offset才能消费到历史数据透传作业则使用latest-offset只消费启动后到达的新消息。八、与课程其他模块的衔接本模块是 Data Engineering Zoomcamp 流处理单元cohorts/2027/07-streaming的扩展实验课程主线使用 Python 消费者与 Kafka 交互而 pyflink 模块把消费端替换为 Flink 分布式引擎展示同一份数据如何被以「声明式 Table API」的方式持续处理。仓库中其他流处理实现如 Python 版 Kafka 消费者与生产者可作为对照阅读理解「手写消费者」与「引擎托管消费窗口计算」两种范式的差异。配套的 homework.md 则提供了基于本模块的进阶练习。参考资料仓库内模块运行指南pyflink/README.md镜像构建定义Dockerfile.flink服务编排配置docker-compose.yml命令封装Makefile作业源码start_job.py、aggregation_job.py、taxi_job.py生产端源码producer.py、load_taxi_data.py【免费下载链接】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),仅供参考