Apache Flink 流式管道实战指南:基于 data-engineer-handbook 的 PyFlink + Kafka + PostgreSQL 全链路搭建
数据工程文档教程【免费下载链接】data-engineer-handbookThis is a repo with links to everything youd 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.envvim flink-env.env环境变量详解example.env 中的完整配置如下KAFKA_WEB_TRAFFIC_SECRETGET FROM WEBSITE KAFKA_WEB_TRAFFIC_KEYGET FROM WEBSITE IP_CODING_KEYMAKE AN ACCOUNT AT https://www.ip2location.io/ TO GET KEY KAFKA_GROUPweb-events KAFKA_TOPICbootcamp-events-prod KAFKA_URLpkc-rgm37.us-west-2.aws.confluent.cloud:9092 FLINK_VERSION1.16.0 PYTHON_VERSION3.7.9 POSTGRES_URLjdbc:postgresql://host.docker.internal:5432/postgres JDBC_BASE_URLjdbc:postgresql://host.docker.internal:5432 POSTGRES_USERpostgres POSTGRES_PASSWORDpostgres POSTGRES_DBpostgres各变量在管道中的实际作用变量说明消费方源码依据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.idKAFKA_TOPIC上游事件主题名create_events_source_kafka 读取该主题KAFKA_URLKafka bootstrap servers 地址所有 Kafka 连接器的properties.bootstrap.serversPOSTGRES_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 的构建逻辑这决定了集群的能力边界基础镜像基于flink:1.16.2注意 README 环境变量中标注的FLINK_VERSION1.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-flink1.16.2、psycopg2-binary2.9.1供脚本内 Python 侧使用、requests供 UDF 调用外部 API安装 Java 11并设置JAVA_HOME下载连接器 JAR到/opt/flink/lib/flink-python-1.16.2.jarPyFlink 运行所需flink-sql-connector-kafka-1.16.2.jarKafka 连接器flink-connector-jdbc-1.16.2.jarJDBC 连接器postgresql-42.2.26.jarPostgreSQL JDBC 驱动。这四类 JAR 是后续CREATE TABLE ... WITH (connector kafka / jdbc)能否执行的物质基础任何缺失都会导致作业在运行时抛出 connector 找不到的异常。集群拓扑docker-compose 中的两个核心服务docker-compose.yml 定义了 Flink 会话集群的最小拓扑jobmanager暴露8081:8081Flink Web UI命令为jobmanager通过FLINK_PROPERTIES设置jobmanager.rpc.address: jobmanager使用extra_hosts: host.docker.internal:host-gateway使容器能访问宿主机上的 PostgreSQLtaskmanager依赖 jobmanager 启动命令带--taskmanager.registration.timeout 5 min设置taskmanager.numberOfTaskSlots: 15与parallelism.default: 3为多并行度聚合作业预留资源。两个服务共用image: eczachly-pyflinkpull_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 已成功注册到 JobManagertaskmanager Successful registration at resource manager akka.tcp://flinkjobmanager: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.pyaggregation_job提交聚合作业运行 aggregation_job.pystop/start停止 / 启动 compose 中的所有服务clean停止并移除容器以及 tag 为none的镜像psql在容器内执行 psql 查询 PostgreSQLpostgres-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_settingssettings)作业以流式模式运行开启每 10 秒的检查点以提供故障恢复能力并行度设为 1。2. 注册自定义 UDFclass 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_typeDataTypes.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 kafkatopic与properties.group.id取自环境变量安全协议SASL_SSLPLAIN机制 JAAS 配置使用KAFKA_WEB_TRAFFIC_KEY/SECRETscan.startup.mode latest-offset与properties.auto.offset.reset latest只消费作业启动后的新事件计算列event_timestamp AS TO_TIMESTAMP(event_time, yyyy-MM-ddTHH:mm:ss.SSSZ)将字符串事件时间解析为时间戳format json。4. 声明 PostgreSQL 汇表create_processed_events_sink_postgres 声明 JDBC sinkconnector 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_hitsTumble.over(lit(5).minutes).on(col(window_timestamp))同时按hostreferrer双维度分组写入第二张聚合表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.modelatest-offset下作业需在事件产生前启动作业崩溃于外部 APIGetLocation的非 200 兜底可防止单条异常记录阻塞管道排查时可先观察 Flink UI 的异常栈与检查点状态。至此你已经掌握了基于本仓库 Flink 训练模块的完整流式管道从环境准备、凭据配置、镜像构建到作业提交、窗口聚合与结果验证。这套模式可直接迁移到你自己的 Kafka Flink PostgreSQL 实时数据处理场景中。赞分享数据工程文档教程【免费下载链接】data-engineer-handbookThis is a repo with links to everything youd ever want to learn about data engineering项目地址https://gitcode.com/GitHub_Trending/da/data-engineer-handbook点击查看免费下载相关推荐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后端大数据流处理批处理上一篇Android Studio 中文语言包安装全记录不写一行代码的全界面汉化方案下一篇免费商用中文字体怎么选思源宋体CN 7个字重从下载到网页上线的完整实操创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关新闻

动手学深度学习:注意力评分函数(掩蔽Softmax、加性注意力与缩放点积注意力)

动手学深度学习:注意力评分函数(掩蔽Softmax、加性注意力与缩放点积注意力)

人工智能深度学习机器学习教程 【免费下载链接】d2l-zh 《动手学深度学习》:面向中文读者、能运行、可讨论。中英文版被70多个国家的500多所大学用于教学。 项目地址: https://gitcode.com/GitHub_Trending/d2/d2l-zh 点击查看 免费下载 导读 在《动手…

2026/10/1 2:34:55 阅读更多 →
HelloGitHub 第 122 期精读:39 个精选开源项目全解析与分类指南

HelloGitHub 第 122 期精读:39 个精选开源项目全解析与分类指南

技术博客文档知识库 【免费下载链接】HelloGitHub :octocat: 分享 GitHub 上有趣、入门级的开源项目。Share interesting, entry-level open source projects on GitHub. 项目地址: https://gitcode.com/GitHub_Trending/he/HelloGitHub 点击查看 免费下载 本篇文章…

2026/10/1 2:33:25 阅读更多 →
type-challenges 第 15 题「最后一个元素 Last of Array」:用类型体操实现数组尾元素提取

type-challenges 第 15 题「最后一个元素 Last of Array」:用类型体操实现数组尾元素提取

示例工程 【免费下载链接】type-challenges Collection of TypeScript type challenges with online judge 项目地址: https://gitcode.com/GitHub_Trending/ty/type-challenges 点击查看 免费下载 type-challenges 是一个带在线判题功能的 TypeScript 类型挑战集合…

2026/10/1 2:32:26 阅读更多 →

最新新闻

Java protected修饰符深度拆解:跨包访问规则与实战避坑指南

Java protected修饰符深度拆解:跨包访问规则与实战避坑指南

刚工作那两年,我一直觉得访问修饰符是Java里“最简单”的知识点——public、private、protected,背一下表格就完事了。直到有次面试官问我:“一个子类在另一个包里,能通过父类的引用访问protected方法吗?”我当时愣了一…

2026/10/1 4:47:10 阅读更多 →
SpringBoot+Vue美发门店管理系统:会员储值与提成设计解析

SpringBoot+Vue美发门店管理系统:会员储值与提成设计解析

做门店管理系统的这些年,我总结出一个规律:凡是老板能坚持用超过三个月的系统,往往不是功能最多的那套,而是把“钱”和“账”算得最清楚的那套。美发门店尤其明显——会员储值、划卡、员工提成,每一笔都牵动着老板的神…

2026/10/1 4:47:10 阅读更多 →
ADE不是IDE升级:智能体运行时的范式迁移

ADE不是IDE升级:智能体运行时的范式迁移

1. 从“写代码的编辑器”到“调度智能体的控制台”:ADE不是IDE的升级版,而是范式迁移的起点你有没有过这种体验:在VS Code里敲完一段Python,运行后发现它只是个“会执行命令的脚本”,而隔壁团队用的那套系统&#xff0…

2026/10/1 4:47:10 阅读更多 →
C语言main函数return值的系统级意义与工程实践

C语言main函数return值的系统级意义与工程实践

1. 为什么“return 0”不是可有可无的装饰,而是程序与操作系统之间的契约刚接触C语言的大一新生,往往在写完第一个printf("hello world!");后,会盯着编辑器里那行return 0;发愣:输出都完成了,它到底在干什么…

2026/10/1 4:47:10 阅读更多 →
Unsloth实战:4-bit量化+LoRA让8G显存轻松微调7B模型

Unsloth实战:4-bit量化+LoRA让8G显存轻松微调7B模型

显存焦虑大概是现在不少想碰微调的人的第一道坎。我自己第一台拿来跑模型的卡只有8G显存,当年跑个7B模型的LoRA,还没等epoch跑完就先收到CUDA Out of Memory,那种卡在进度条上的感觉,谁遇过谁知道。后来试到Unsloth,同…

2026/10/1 4:47:09 阅读更多 →
科研第二曲线:实验卡壳时如何继续推进课题进度

科研第二曲线:实验卡壳时如何继续推进课题进度

“课题进度总卡在实验上?你可能忽略了科研的‘第二曲线’”说实话,我见过太多人把“科研进度”和“实验进度”直接画等号了。实验一不顺利,整个人就进入“等待模式”:做不出来就干等着,等仪器、等样品、等表征、等结果…

2026/10/1 4:46:09 阅读更多 →

日新闻

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/1 0:00:30 阅读更多 →
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/1 0:00:30 阅读更多 →
黑夜航拍船只数据集训练YOLOV5模型全流程解析

黑夜航拍船只数据集训练YOLOV5模型全流程解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/1 1:01:17 阅读更多 →

周新闻

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解

如何划分训练/验证集:Spirula Studio五种eval_mode策略详解 【免费下载链接】spirula-studio Cross-vendor 3D Gaussian Splatting trainer - video to splat to mesh, Vulkan or CUDA. 项目地址: https://gitcode.com/GitHub_Trending/sp/spirula-studio Sp…

2026/9/30 13:14:22 阅读更多 →
SEO怎么推广速查手册新手避坑实战指南

SEO怎么推广速查手册新手避坑实战指南

SEO怎么推广速查手册新手避坑实战指南 模板网站太丑不够用?别急着加滤镜,那是治标不治本。很多老板盯着后台流量掉得眼红,却还在纠结首页Banner的圆角是不是3像素。这就像穿着西装去挖土,姿势不对,努力白费。我整理这份 速查手册…

2026/9/30 18:13:06 阅读更多 →
FireRed-OpenStoryline少样本仿写深度解析:AI Agent如何复刻你的独特文案风格与节奏

FireRed-OpenStoryline少样本仿写深度解析:AI Agent如何复刻你的独特文案风格与节奏

FireRed-OpenStoryline少样本仿写深度解析:AI Agent如何复刻你的独特文案风格与节奏 【免费下载链接】FireRed-OpenStoryline FireRed-OpenStoryline is an AI video editing agent that transforms manual editing into intention-driven directing through natural language …

2026/9/30 13:14:49 阅读更多 →

月新闻

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

我发现了一个新思路:用 Remotion + Claude Code 像写代码一样自动化生成短视频

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/1 0:00:30 阅读更多 →
Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

Windows下 Codex 中 Chrome 和 Computer Use 插件不可用问题排查及解决参考方式:TaoToken 统一 Key 配置与验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/1 0:00:30 阅读更多 →
黑夜航拍船只数据集训练YOLOV5模型全流程解析

黑夜航拍船只数据集训练YOLOV5模型全流程解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

2026/10/1 1:01:17 阅读更多 →