- 数据工程
- 文档
- 教程
【免费下载链接】data-engineer-handbook
This is a repo with links to everything you'd ever want to learn about data engineering
本篇技术指南以># On Ubuntu or Debian: sudo apt-get update sudo apt-get install build-essential # On CentOS or Fedora: sudo dnf install make # On macOS: xcode-select --install # On Windows: choco install make # uses Chocolatey
如果不想安装 Make,也完全可以直接复制 Makefile 中对应 target 后的命令,在终端手动执行——本文后续每个步骤都会同时给出两种方式。
获取代码并进入模块目录
git clone https://gitcode.com/GitHub_Trending/da/data-engineer-handbook.git cd intermediate-bootcamp/materials/4-apache-flink-training配置凭据:从 example.env 到 flink-env.env
复制环境文件
模块通过环境变量注入 Kafka 凭据与数据库连接信息。第一步是把模板文件复制为实际使用的环境文件:
cp example.env flink-env.env然后用vim或任意编辑器修改flink-env.env:
vim flink-env.env环境变量详解
example.env 中的完整配置如下:
KAFKA_WEB_TRAFFIC_SECRET="<GET FROM WEBSITE>" KAFKA_WEB_TRAFFIC_KEY="<GET FROM WEBSITE>" IP_CODING_KEY="MAKE AN ACCOUNT AT https://www.ip2location.io/ TO GET KEY" KAFKA_GROUP=web-events KAFKA_TOPIC=bootcamp-events-prod KAFKA_URL=pkc-rgm37.us-west-2.aws.confluent.cloud:9092 FLINK_VERSION=1.16.0 PYTHON_VERSION=3.7.9 POSTGRES_URL=jdbc:postgresql://host.docker.internal:5432/postgres JDBC_BASE_URL=jdbc:postgresql://host.docker.internal:5432 POSTGRES_USER=postgres POSTGRES_PASSWORD=postgres POSTGRES_DB=postgres各变量在管道中的实际作用:
| 变量 | 说明 | 消费方(源码依据) |
|---|---|---|
KAFKA_WEB_TRAFFIC_KEY/KAFKA_WEB_TRAFFIC_SECRET | Confluent Cloud Kafka 的 SASL 认证凭据 | start_job.py 用于构造PlainLoginModule的 JAAS 配置 |
IP_CODING_KEY | ip2location.io 地理定位 API 密钥 | start_job.py 中GetLocationUDF 调用 API 时传入 |
KAFKA_GROUP | Kafka 消费者组 ID | 源表/汇表 DDL 中的properties.group.id |
KAFKA_TOPIC | 上游事件主题名 | create_events_source_kafka 读取该主题 |
KAFKA_URL | Kafka bootstrap servers 地址 | 所有 Kafka 连接器的properties.bootstrap.servers |
POSTGRES_URL/JDBC_BASE_URL | JDBC 连接串,指向宿主机上的 PostgreSQL | JDBC sink 的url配置 |
POSTGRES_USER/POSTGRES_PASSWORD/POSTGRES_DB | PostgreSQL 连接凭据 | 同时被容器环境变量与 JDBC sink 使用 |
安全警告:flink-env.env中保存的是云上 Kafka 资源的真实凭据,严禁将其推送或分享到训练营之外,否则可能导致云端资源被他人滥用。其余关于凭据的旧版说明可以忽略——仓库更新后,需要的一切都已包含在example.env中。如果修改了容器化 PostgreSQL 的POSTGRES_USER与POSTGRES_PASSWORD,请保持环境文件与 docker-compose.yml 中的默认值一致,否则连接会失败;不修改则保持postgres/postgres默认值即可。
理解 Flink 镜像构建:Dockerfile 逐层拆解
在运行管道前,先理解 Dockerfile 的构建逻辑,这决定了集群的能力边界:
- 基础镜像:基于
flink:1.16.2(注意 README 环境变量中标注的FLINK_VERSION=1.16.0与镜像实际使用的 1.16.2 存在小版本差异,以镜像构建产物为准); - 安装 Python 3.7.9:官方 PyFlink 当时仅正式支持 Python 3.6/3.7/3.8,而 Debian 11 自带 Python 3.9,因此需要从源码编译 3.7.9(步骤:下载源码 →
./configure --enable-shared→make→make install→ 软链python); - 安装 PyFlink 依赖:通过 requirements.txt 安装
apache-flink==1.16.2、psycopg2-binary==2.9.1(供脚本内 Python 侧使用)、requests(供 UDF 调用外部 API); - 安装 Java 11并设置
JAVA_HOME; - 下载连接器 JAR到
/opt/flink/lib/:flink-python-1.16.2.jar:PyFlink 运行所需;flink-sql-connector-kafka-1.16.2.jar:Kafka 连接器;flink-connector-jdbc-1.16.2.jar:JDBC 连接器;postgresql-42.2.26.jar:PostgreSQL JDBC 驱动。
这四类 JAR 是后续CREATE TABLE ... WITH ('connector' = 'kafka' / 'jdbc')能否执行的物质基础,任何缺失都会导致作业在运行时抛出 connector 找不到的异常。
集群拓扑:docker-compose 中的两个核心服务
docker-compose.yml 定义了 Flink 会话集群的最小拓扑:
- jobmanager:暴露
8081:8081(Flink Web UI),命令为jobmanager;通过FLINK_PROPERTIES设置jobmanager.rpc.address: jobmanager;使用extra_hosts: host.docker.internal:host-gateway使容器能访问宿主机上的 PostgreSQL; - taskmanager:依赖 jobmanager 启动,命令带
--taskmanager.registration.timeout 5 min;设置taskmanager.numberOfTaskSlots: 15与parallelism.default: 3,为多并行度聚合作业预留资源。
两个服务共用image: eczachly-pyflink(pull_policy: never,必须由本地build生成),并共享卷挂载./src/:/opt/src,使作业脚本能被 JobManager 读取。PostgreSQL 不在本 compose 文件中,需要单独通过 Makefile 的db-init或训练营前几周的容器启动。
运行管道:从构建到验证的完整流程
第 1 步:启动 Flink 集群
make up # 没有 make 时手动执行: docker compose --env-file flink-env.env up --build --remove-orphans -d该命令会构建基础镜像并启动 Flink 集群(注意:make up本身并不包含 PostgreSQL——PostgreSQL 需按前几周教程另行启动,或用make db-init)。首次构建镜像需要 5 到 30 分钟,后续重建只需几秒(只要没有删除镜像)。
务必等待 Flink Web UI 就绪(访问 http://localhost:8081/)再进入下一步。镜像构建完成后,Docker 会自动拉起 jobmanager 与 taskmanager 服务,大约需要一分钟。观察容器日志,当出现以下日志行时,说明 TaskManager 已成功注册到 JobManager:
taskmanager Successful registration at resource manager akka.tcp://flink@jobmanager:6123/user/rpc/resourcemanager_* under registration id <id_number>第 2 步:初始化 PostgreSQL 目标表
在本地(或容器化)PostgreSQL 上执行 sql/init.sql,创建下游 sink 表:
CREATE TABLE IF NOT EXISTS processed_events ( ip VARCHAR, event_timestamp TIMESTAMP(3), referrer VARCHAR, host VARCHAR, url VARCHAR, geodata VARCHAR );该表结构与 PyFlink 作业中create_processed_events_sink_postgres定义的 JDBC sink 表字段一一对应,是作业能成功写入的前提。
第 3 步:提交 PyFlink 作业
make job # 没有 make 时手动执行: docker compose exec jobmanager ./bin/flink run -py /opt/src/job/start_job.py --pyFiles /opt/src -d大约一分钟后,终端会提示作业提交成功(例如Job has been submitted with JobID <job_id_number>)。回到 Flink Web UI 的 http://localhost:8081/#/job/running 页面即可看到作业正在运行。
第 4 步:触发事件并验证落库
访问训练营提供的事件触发页面,即可向上游 Kafka 主题产生一条新的 Web 事件。随后查询 PostgreSQL 确认数据已写入:
make psql其实际执行的是进入容器并连接数据库:
docker exec -it eczachly-flink-postgres psql -U postgres -d postgres在 psql 中验证:
postgres=# SELECT COUNT(*) FROM processed_events; count ------- 739 (1 row)只要计数在增长,就证明整条 Kafka → Flink → PostgreSQL 管道已打通。
第 5 步:停止与清理
make stop # 停止正在运行的 compose 服务 make down # 停止并移除 compose 服务 make clean # 移除容器与 <none> 悬空镜像数据持久化说明:PostgreSQL 容器内的/var/lib/postgresql/data挂载到了本机./postgres-data目录,因此即使停止或移除容器,容器内写入的数据也不会丢失。
Makefile 全量命令速查
在模块目录运行make help可随时查看所有可用命令,当前支持的命令如下:
| 命令 | 作用 |
|---|---|
help | 显示帮助 |
db-init | 构建并运行 PostgreSQL 数据库服务 |
build | 构建带 PyFlink 与连接器的 Flink 基础镜像 |
up | 构建基础镜像并启动 Flink 集群 |
down | 关闭 Flink 集群 |
job | 提交 Flink 作业(运行 start_job.py) |
aggregation_job | 提交聚合作业(运行 aggregation_job.py) |
stop/start | 停止 / 启动 compose 中的所有服务 |
clean | 停止并移除容器以及 tag 为<none>的镜像 |
psql | 在容器内执行 psql 查询 PostgreSQL |
postgres-die-mac/postgres-die-pc | 删除本机(Mac / PC)与 Docker 中挂载的 postgres 数据目录 |
源码深潜:start_job.py 的流式管道实现
src/job/start_job.py 是本模块的核心作业,其执行流程如下:
1. 环境初始化与检查点
env = StreamExecutionEnvironment.get_execution_environment() env.enable_checkpointing(10 * 1000) # 每 10 秒一次检查点 env.set_parallelism(1) settings = EnvironmentSettings.new_instance().in_streaming_mode().build() t_env = StreamTableEnvironment.create(env, environment_settings=settings)作业以流式模式运行,开启每 10 秒的检查点以提供故障恢复能力,并行度设为 1。
2. 注册自定义 UDF
class GetLocation(ScalarFunction): def eval(self, ip_address): response = requests.get("https://api.ip2location.io", params={ 'ip': ip_address, 'key': os.environ.get("IP_CODING_KEY") }) data = json.loads(response.text) return json.dumps({'country': data.get('country_code', ''), 'state': data.get('region_name', ''), 'city': data.get('city_name', '')}) get_location = udf(GetLocation(), result_type=DataTypes.STRING()) t_env.create_temporary_function("get_location", get_location)GetLocation继承ScalarFunction,对每条记录的 IP 发起 HTTP 请求,返回国家、州、城市的 JSON 字符串;请求失败时返回空对象{},避免单条坏数据中断整个作业。调用失败时返回空 dict 的兜底逻辑(response.status_code != 200)体现了流式作业对上游异常的可恢复性设计。
3. 声明 Kafka 源表
create_events_source_kafka 通过 Flink SQL DDL 声明源表,关键配置包括:
connector = 'kafka',topic与properties.group.id取自环境变量;- 安全协议
SASL_SSL+PLAIN机制 + JAAS 配置(使用KAFKA_WEB_TRAFFIC_KEY/SECRET); scan.startup.mode = 'latest-offset'与properties.auto.offset.reset = 'latest':只消费作业启动后的新事件;- 计算列:
event_timestamp AS TO_TIMESTAMP(event_time, 'yyyy-MM-dd''T''HH:mm:ss.SSS''Z'''),将字符串事件时间解析为时间戳; format = 'json'。
4. 声明 PostgreSQL 汇表
create_processed_events_sink_postgres 声明 JDBC sink:
'connector' = 'jdbc', 'url' = os.environ.get("POSTGRES_URL"), 'table-name' = 'processed_events', 'username' / 'password' 从环境变量读取, 'driver' = 'org.postgresql.Driver'注意url使用host.docker.internal——这正是 docker-compose.yml 中extra_hosts映射宿主机网关的用武之地,使容器内作业可以访问宿主机上的 PostgreSQL。
5. 组装 INSERT 查询
t_env.execute_sql(f""" INSERT INTO {postgres_sink} SELECT ip, event_timestamp, referrer, host, url, get_location(ip) as geodata FROM {source_table} """).wait()每条 Kafka 事件经get_location(ip)增强后写入processed_events,.wait()阻塞至作业完成提交。
进阶扩展:aggregation_job.py 的窗口聚合
src/job/aggregation_job.py 展示了在流上做滚动窗口(Tumbling Window)聚合的写法:
- 源表通过计算列
window_timestamp AS TO_TIMESTAMP(event_time, ...)解析事件时间,并用WATERMARK FOR window_timestamp AS window_timestamp - INTERVAL '15' SECOND声明 15 秒的水位线,容忍乱序数据; - 作业以并行度 3 运行(
env.set_parallelism(3)),对应 taskmanager 的parallelism.default: 3; - 对每个 5 分钟窗口按
host分组统计num_hits(Tumble.over(lit(5).minutes).on(col("window_timestamp"))),同时按host+referrer双维度分组写入第二张聚合表processed_events_aggregated_source; - 结果通过 JDBC 写入 PostgreSQL。
make aggregation_job即可提交该作业。这一示例揭示了从「原始事件管道」到「指标聚合管道」的演进路径,也是本模块作业的核心素材——按 IP 与 host 做 5 分钟 gap 的会话化(sessionization),并回答「Tech Creator 上单个用户会话的平均事件数」等问题。
验证与故障排查要点
- UI 未就绪:检查
docker compose logs jobmanager,确认 TaskManager 注册日志出现后再提交作业; - Kafka 认证失败:核对
flink-env.env中KAFKA_WEB_TRAFFIC_KEY/SECRET与 JAAS 格式(注意转义引号); - 写入失败:确认已执行 sql/init.sql 且 PostgreSQL 凭据与 compose 环境变量一致;
- 无数据消费:确认上游事件已触发,且
scan.startup.mode=latest-offset下作业需在事件产生前启动; - 作业崩溃于外部 API:
GetLocation的非 200 兜底可防止单条异常记录阻塞管道,排查时可先观察 Flink UI 的异常栈与检查点状态。
至此,你已经掌握了基于本仓库 Flink 训练模块的完整流式管道:从环境准备、凭据配置、镜像构建,到作业提交、窗口聚合与结果验证。这套模式可直接迁移到你自己的 Kafka + Flink + PostgreSQL 实时数据处理场景中。
- 数据工程
- 文档
- 教程
【免费下载链接】data-engineer-handbook
This is a repo with links to everything you'd ever want to learn about data engineering
相关推荐
Data Engineering Zoomcamp:用 Apache Flink 构建端到端 PyFlink 流式管道实战指南
Data Engineering Zoomcamp:用 Apache Flink 构建端到端 PyFlink 流式管道实战指南 Apache Flink 是当前
教程数据工程GB28181开源视频监控平台:多品牌摄像头一套平台搞定,免费开箱即用,十分钟能看到画面吗?
GB28181开源视频监控平台:多品牌摄像头一套平台搞定,免费开箱即用,十分钟能看到画面吗? wvp GB28181 pro 是一款免费开源、可商用的 GB28
后端音视频前端Apache Flink PyFlink 指南:用 Python API 构建批流一体的数据管道
Apache Flink PyFlink 指南:用 Python API 构建批流一体的数据管道 PyFlink 是 Apache Flink 的 Python
后端大数据流处理批处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考