拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Flink 流式管道实战指南:基于 data-engineer-handbook 的 PyFlink + Kafka + PostgreSQL 全链路搭建

Apache Flink 流式管道实战指南:基于 data-engineer-handbook 的 PyFlink + Kafka + PostgreSQL 全链路搭建
  • 数据工程
  • 文档
  • 教程

【免费下载链接】data-engineer-handbook

This is a repo with links to everything you'd ever want to learn about data engineering

项目地址:https://gitcode.com/GitHub_Trending/da/data-engineer-handbook
点击查看免费下载

本篇技术指南以># 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_SECRETConfluent Cloud Kafka 的 SASL 认证凭据start_job.py 用于构造PlainLoginModule的 JAAS 配置
IP_CODING_KEYip2location.io 地理定位 API 密钥start_job.py 中GetLocationUDF 调用 API 时传入
KAFKA_GROUPKafka 消费者组 ID源表/汇表 DDL 中的properties.group.id
KAFKA_TOPIC上游事件主题名create_events_source_kafka 读取该主题
KAFKA_URLKafka bootstrap servers 地址所有 Kafka 连接器的properties.bootstrap.servers
POSTGRES_URL/JDBC_BASE_URLJDBC 连接串,指向宿主机上的 PostgreSQLJDBC sink 的url配置
POSTGRES_USER/POSTGRES_PASSWORD/POSTGRES_DBPostgreSQL 连接凭据同时被容器环境变量与 JDBC sink 使用

安全警告:flink-env.env中保存的是云上 Kafka 资源的真实凭据,严禁将其推送或分享到训练营之外,否则可能导致云端资源被他人滥用。其余关于凭据的旧版说明可以忽略——仓库更新后,需要的一切都已包含在example.env中。如果修改了容器化 PostgreSQL 的POSTGRES_USER与POSTGRES_PASSWORD,请保持环境文件与 docker-compose.yml 中的默认值一致,否则连接会失败;不修改则保持postgres/postgres默认值即可。

理解 Flink 镜像构建:Dockerfile 逐层拆解

在运行管道前,先理解 Dockerfile 的构建逻辑,这决定了集群的能力边界:

  1. 基础镜像:基于flink:1.16.2(注意 README 环境变量中标注的FLINK_VERSION=1.16.0与镜像实际使用的 1.16.2 存在小版本差异,以镜像构建产物为准);
  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);
  3. 安装 PyFlink 依赖:通过 requirements.txt 安装apache-flink==1.16.2、psycopg2-binary==2.9.1(供脚本内 Python 侧使用)、requests(供 UDF 调用外部 API);
  4. 安装 Java 11并设置JAVA_HOME;
  5. 下载连接器 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

项目地址:https://gitcode.com/GitHub_Trending/da/data-engineer-handbook
点击查看免费下载
上一篇:Android Studio 中文语言包安装全记录:不写一行代码的全界面汉化方案
下一篇:免费商用中文字体怎么选?思源宋体CN 7个字重从下载到网页上线的完整实操

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

返回列表