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-1 | redpandadata/redpanda:v24.2.18 | 9092 / 29092 / 8082 / 28082 | Kafka 协议兼容的消息队列,承担 Kafka 角色 |
jobmanager | pyflink:1.16.0(本地构建) | 8081(Flink UI) | Flink 作业管理与调度,提交 PyFlink 作业的入口 |
taskmanager | pyflink:1.16.0 | 6121 / 6122 | 执行流处理算子,默认 15 个任务槽、并行度 3 |
postgres | postgres:14 | 5432 | 流处理结果的落地目标 |
数据流方向为:生产者把消息写入 Redpanda topic → PyFlink 作业从 topic 消费 → 经流式处理(透传或窗口聚合)→ 写入 PostgreSQL 表。
二、环境准备:Docker、Docker Compose 与 Make
运行本模块需要三个前置组件(对应 README.md 中的 Installation 部分):
- Docker(必需)——承载 Flink 集群、Redpanda 与 PostgreSQL;
- Docker 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/pyflink):
cd cohorts/2027/07-streaming/extras/pyflink三、一键拉起:构建镜像并启动 Flink 集群
3.1 make up 做了什么
在模块根目录执行:
make up等价于手动执行(见 Makefile):
docker 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 8; - Python 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-flink==1.16.0、psycopg2-binary==2.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—— Kafka(Redpanda)连接器;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://flink@jobmanager:6123/user/rpc/resourcemanager_* under registration id <id_number>然后访问 Flink Web UI:http://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_servers='localhost:9092'能连通的原因; - JobManager 环境变量:通过
POSTGRES_URL、POSTGRES_USER、POSTGRES_PASSWORD、POSTGRES_DB注入 PostgreSQL 连接信息,且带默认值postgres,可用.env文件覆盖;FLINK_PROPERTIES中声明jobmanager.rpc.address: jobmanager; - TaskManager 并行配置:
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等价于手动执行(Makefile):
docker compose exec jobmanager ./bin/flink run -py /opt/src/job/start_job.py --pyFiles /opt/src -d命令拆解:
docker compose exec jobmanager:进入 JobManager 容器执行命令;./bin/flink run:Flink 命令行提交作业;-py /opt/src/job/start_job.py:指定要运行的 Python 作业文件(挂载自仓库 src/job/start_job.py);--pyFiles /opt/src:把整个src目录加入 Python 依赖路径,供作业内 import 使用;-d:detached 模式,提交后立即返回,作业在后台运行。
大约一分钟后,终端会出现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 数据源(topic:test-topic):
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' );PostgreSQL 汇入目标(表名:processed_events):
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' );作业通过一条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 ENFORCED(JDBC 汇入目标支持按主键 upsert)。
运行该作业前需要先启动生产者,并且聚合作业消费的是earliest偏移,因此会把历史上已发送到test-topic的消息一并按分钟窗口统计,得到类似「某分钟窗口内某个 test_data 值出现了多少次」的聚合结果。
4.5 出租车真实数据作业
模块还附带了一个面向真实业务数据的作业 taxi_job.py:
- 数据源:Kafka topic
green-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 集群(含 PostgreSQL) |
make 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:9092(Redpanda 的 OUTSIDE 监听),而容器内 Flink 作业使用redpanda-1:29092(PLAINTEXT 监听),两者不能混用。 - 聚合作业没有数据:确认生产者已先启动并向
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),仅供参考