Data Engineering Zoomcamp:使用 PyFlink 搭建 Kafka 到 PostgreSQL 的实时流处理管道
2026/9/12 4:56:14 网站建设 项目流程

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(本地构建)8081(Flink UI)Flink 作业管理与调度,提交 PyFlink 作业的入口
taskmanagerpyflink:1.16.06121 / 6122执行流处理算子,默认 15 个任务槽、并行度 3
postgrespostgres:145432流处理结果的落地目标

数据流方向为:生产者把消息写入 Redpanda topic → PyFlink 作业从 topic 消费 → 经流式处理(透传或窗口聚合)→ 写入 PostgreSQL 表。

二、环境准备:Docker、Docker Compose 与 Make

运行本模块需要三个前置组件(对应 README.md 中的 Installation 部分):

  1. Docker(必需)——承载 Flink 集群、Redpanda 与 PostgreSQL;
  2. Docker Compose(必需)——用于一次拉起全部服务;
  3. 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

这一条命令完成三件事:

  1. 依据 Dockerfile.flink 构建名为pyflink:1.16.0的 Flink 基础镜像;
  2. 依次启动 Redpanda、JobManager、TaskManager、PostgreSQL 四个服务;
  3. 创建 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.0psycopg2-binary==2.9.1requestskafka-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.pybootstrap_servers='localhost:9092'能连通的原因;
  • JobManager 环境变量:通过POSTGRES_URLPOSTGRES_USERPOSTGRES_PASSWORDPOSTGRES_DB注入 PostgreSQL 连接信息,且带默认值postgres,可用.env文件覆盖;FLINK_PROPERTIES中声明jobmanager.rpc.address: jobmanager
  • TaskManager 并行配置taskmanager.numberOfTaskSlots: 15parallelism.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.py

producer.py 的核心逻辑是:用kafka-pythonKafkaProducer连接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_hourtest_datanum_hits三个字段,并声明PRIMARY KEY (event_hour, test_data) NOT ENFORCED(JDBC 汇入目标支持按主键 upsert)。

运行该作业前需要先启动生产者,并且聚合作业消费的是earliest偏移,因此会把历史上已发送到test-topic的消息一并按分钟窗口统计,得到类似「某分钟窗口内某个 test_data 值出现了多少次」的聚合结果。

4.5 出租车真实数据作业

模块还附带了一个面向真实业务数据的作业 taxi_job.py:

  • 数据源:Kafka topicgreen-data,消费模式earliest-offset,格式为 JSON,包含 20 个出租车行程字段(VendorIDlpep_pickup_datetimetrip_distancetotal_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-macmake 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询