data-engineering-zoomcamp 07-streaming 实战:用 PyFlink 编写首个 Pass-through 作业,把 Kafka 事件写入 PostgreSQL

发布时间:2026/9/12 5:59:56

data-engineering-zoomcamp 07-streaming 实战:用 PyFlink 编写首个 Pass-through 作业,把 Kafka 事件写入 PostgreSQL data-engineering-zoomcamp 07-streaming 实战用 PyFlink 编写首个 Pass-through 作业把 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本指南聚焦>def create_events_source_kafka(t_env): table_name events source_ddl f CREATE TABLE {table_name} ( PULocationID INTEGER, DOLocationID INTEGER, trip_distance DOUBLE, total_amount DOUBLE, tpep_pickup_datetime BIGINT ) WITH ( connector kafka, properties.bootstrap.servers redpanda:29092, topic rides, scan.startup.mode latest-offset, properties.auto.offset.reset latest, format json ); t_env.execute_sql(source_ddl) return table_name逐项拆解列定义与 producer 的 JSON 字段一一对应PULocationID、DOLocationID、trip_distance、total_amount、tpep_pickup_datetime。这些字段来自 producer.py 发送的消息而消息结构定义在 models.py 的Ridedataclass 中——注意tpep_pickup_datetime的类型是intepoch 毫秒所以源表中对应声明为BIGINTproperties.bootstrap.servers redpanda:29092这是 Docker 内部网络地址。Flink 运行在容器里不能像本机脚本那样用localhost:9092——localhost会指向 Flink 容器自身。29092是 Redpanda 暴露在 Docker 网络内部的 PLAINTEXT 监听端口见 docker-compose.yml 中 Redpanda 的--advertise-kafka-addr配置scan.startup.mode latest-offset作业启动后只读取新到达的消息不重放历史数据。三种可选模式的对比latest-offset/earliest-offset/timestamp在 09-offsets-earliest-vs-latest.md 中有完整表格说明properties.auto.offset.reset latest与scan.startup.mode配套的 Kafka 消费者参数同样控制起始读取位置format jsonFlink 自动完成 JSON 反序列化无需像 Python 消费者那样手写ride_deserializer。这段 DDL 通过t_env.execute_sql(source_ddl)注册到 Table 环境返回表名events供后续 SQL 引用。第二步声明 PostgreSQL 目标表第二块拼图是 PostgreSQL 的目标表同样只是一段声明式 DDLdef create_processed_events_sink_postgres(t_env): table_name processed_events sink_ddl f CREATE TABLE {table_name} ( PULocationID INTEGER, DOLocationID INTEGER, trip_distance DOUBLE, total_amount DOUBLE, pickup_datetime TIMESTAMP ) WITH ( connector jdbc, url jdbc:postgresql://postgres:5432/postgres, table-name {table_name}, username postgres, password postgres, driver org.postgresql.Driver ); t_env.execute_sql(sink_ddl) return table_name关键点connector jdbc启用 JDBC 连接器底层依赖 Dockerfile.flink 中下载的flink-connector-jdbc-core、flink-connector-jdbc-postgres与 PostgreSQL 驱动 JARurl jdbc:postgresql://postgres:5432/postgres同样使用 Docker 内部服务名postgres而非localhosttable-name processed_events对应 05-save-events-to-postgresql.md 中创建的物理表字段完全一致唯一差异是 Flink 端pickup_datetime声明为TIMESTAMP类型driver org.postgresql.Driver显式指定 JDBC 驱动类。注意这里没有 psycopg2、没有 INSERT 语句——只需声明表结构Flink 负责其余一切建立连接、按 schema 反序列化、批量写入、以及失败时的重试与一致性保障。第三步组装执行管线最后一步是把两条 DDL 与一条流式 SQL 查询组装成完整作业from pyflink.datastream import StreamExecutionEnvironment from pyflink.table import EnvironmentSettings, StreamTableEnvironment def log_processing(): env StreamExecutionEnvironment.get_execution_environment() env.enable_checkpointing(10 * 1000) # checkpoint every 10 seconds settings EnvironmentSettings.new_instance().in_streaming_mode().build() t_env StreamTableEnvironment.create(env, environment_settingssettings) source_table create_events_source_kafka(t_env) postgres_sink create_processed_events_sink_postgres(t_env) t_env.execute_sql( f INSERT INTO {postgres_sink} SELECT PULocationID, DOLocationID, trip_distance, total_amount, TO_TIMESTAMP_LTZ(tpep_pickup_datetime, 3) as pickup_datetime FROM {source_table} ).wait() if __name__ __main__: log_processing()这段代码揭示了 Flink 的编程模型Streaming modeEnvironmentSettings.new_instance().in_streaming_mode().build()声明作业处于流式模式——作业持续运行、等待新数据而不是处理完就退出INSERT INTO ... SELECT即管线本身从events源表读取 → 将tpep_pickup_datetimeBIGINT 毫秒时间戳用TO_TIMESTAMP_LTZ(..., 3)转换为带毫秒精度precision 3的TIMESTAMP→ 写入processed_events目标表。这一步取代了 Python 版本中的datetime.fromtimestamp(... / 1000)手工转换StreamExecutionEnvironment与StreamTableEnvironment的关系前者是 DataStream API 的执行环境后者基于它构建 Table API 环境从而能执行 SQL DDL 与查询.wait()阻塞等待作业提交并持续运行期间 Flink 会不断消费事件、执行转换、写入 sink。仓库中的 pass_through_job.py 实际实现还包了一层try/except在异常时打印 Writing records from Kafka to JDBC failed 的调试信息可作为生产代码的参考骨架。Checkpoint10 秒一次的状态快照为什么重要env.enable_checkpointing(10 * 1000)是本课的理论核心之一。它告诉 Flink 每 10 秒对作业状态做一次快照。一次 checkpoint 会捕获Kafka 偏移量Flink 已读到主题的哪个位置在途数据正处于处理过程中、尚未落库的事件窗口状态在有窗口聚合的作业中尚未关闭的窗口内部状态。如果作业崩溃它会从最近一次 checkpoint 恢复而不是从头重新消费。这对带窗口的作业尤为关键假设你有一个 5 分钟窗口作业运行到第 2 分钟时失败——Flink 不仅记录偏移量还会把已经填充了 2 分钟数据的窗口序列化到磁盘重启后直接从断点续跑半满的窗口原样恢复不会丢数据也不会重复计算。代价是韧性 vs 效率的权衡每 1 秒 checkpoint 一次Flink 必须频繁地序列化并持久化全部状态开销很大每 10 分钟 checkpoint 一次故障时最多丢失 10 分钟的处理进度需要重放这段时间的消息10 秒对大多数作业是合理的默认值——兼顾恢复粒度与性能开销。需要强调一个常见误区详见 09-offsets-earliest-vs-latest.mdcheckpoint 归属于特定的作业实例。如果你在 Flink UI 上取消作业再重新提交这是一个全新的作业它不认识旧的 checkpoint。此时scan.startup.mode会重新决定读取起点——earliest-offset会从头重放整个主题可能产生重复数据latest-offset则只收新消息。offset 设置只在作业启动时生效作业运行后是 checkpoint 在接管进度跟踪。提交作业并验证端到端链路1. 提交作业docker compose exec jobmanager ./bin/flink run \ -py /opt/src/job/pass_through_job.py \ --pyFiles /opt/src -d命令解读docker compose exec jobmanager进入 jobmanager 容器执行命令因为 Flink CLI 位于容器内./bin/flink runFlink 的作业提交入口-py /opt/src/job/pass_through_job.py指定 Python 作业文件./src/已被挂载到容器的/opt/src见 docker-compose.yml--pyFiles /opt/src把源码目录加入 Python 模块搜索路径使作业内可以 import 其他模块-ddetached 模式命令立即返回作业在后台持续运行。成功后输出类似Job has been submitted with JobID 663cff6811b65e97fc1e068d641401f42. 查看 Flink UI打开 http://localhost:8081对应 jobmanager 的 Web UI 端口可以看到一个 Running 状态的作业。由于源表配置了latest-offset此时作业正在等待新消息。3. 发送数据uv run python src/producers/producer.pyproducer.py 会下载 NYC 黄色出租车 2025 年 11 月数据的前 1000 行逐条序列化为 JSON 发送到rides主题。4. 验证落库SELECT count(*) FROM processed_events;对比之前的 Python 消费者方案结果相同但 checkpoint、offset 管理与 PostgreSQL 写入全部由 Flink 自动完成。查询后如需清理可参考 13-cleanup.mddocker compose down停掉服务加-v参数同时删除持久化数据卷。小结本课要点与下一步pass-through 作业是 Flink 流处理的最小完整范式它奠定了后续所有进阶作业的骨架声明式优于命令式用CREATE TABLE ... WITH (...)声明源与目标用INSERT INTO ... SELECT描述数据流Flink 自动生成执行计划并分布式运行容器内寻址Flink 在 Docker 网络中Kafka 与 PostgreSQL 都必须用服务名而非localhostcheckpoint 是故障恢复的基石它同时覆盖偏移量与状态选择间隔需要在恢复粒度与开销之间权衡offset 只在启动时生效latest-offset收新消息、earliest-offset重放历史、timestamp从指定时间点恢复。掌握了透传作业后下一步就是本模块的进阶主题在 10-aggregation-with-tumbling-windows.md 中用 tumbling window 做时间窗口聚合届时你会看到 checkpoint 对窗口状态的保护正是保障聚合结果不重不丢的关键aggregation_job.py的实现同样位于 code/src/job/ 目录可作为对照学习。【免费下载链接】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),仅供参考
延伸阅读

更多相关文章

2026/9/12 6:50:01

macOS 应用精选集 awesome-macOS:告别盲目找软件的烦恼

macOS 应用精选集 awesome-macOS:告别盲目找软件的烦恼 【免费下载链接】awesome-macOS  A curated list of awesome applications, softwares, tools and shiny things for macOS. 项目地址: https://gitcode.com/GitHub_Trending/aw/awesome-macOS 装个…

2026/9/12 6:50:01

Linux wheel组与sudo权限完全指南:从原理到安全配置

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

2026/9/12 6:50:01

Virtex-7 FPGA选型与工程实践:资源、应用与调试要点

做项目选型时拿到过Virtex-7的样片,也帮人排查过这块芯片在高速数据采集板卡上的各种问题。说实话,Xilinx 7系列在FPGA圈子里地位很特殊——从2010年左右发布到现在十几年了,依然是很多军工、通信、测试测量项目的首选主力。而Virtex-7作为7系…

2026/9/12 6:45:00

C++20策略内联与std::ranges性能优化解析

1. 理解std::ranges与策略内联的本质当我在2019年首次接触C20的ranges库时,最让我震撼的不是它的管道操作符语法糖,而是隐藏在背后的编译期魔法。策略内联编译器(Policy-Based Inlining Compiler)正是这种魔法的核心引擎&#xff…

2026/9/12 2:05:33

超人会飞不算本事:系统稳定依赖清晰规则与边界设计

开头先不绕弯子。“#斯坦李吐槽dc 所以超人是无缘无故会飞的嘛哈哈哈哈哈哈哈锤哥真是技术人才啊!#雷神 #复联”这类调侃式短标题,第一波冲击力在于它把两个宇宙的角色塞进同一个吐槽箱里,但细想一下就能发现,它真正碰到的根本不是…

2026/9/12 3:55:12

超人VS蜘蛛侠:拆解超级IP的影响力与传播方法论

把“蜘蛛侠 vs 超人”放在 CSDN 上聊,可能很多人第一反应是走错片场了。但如果把这两个角色看成“两个持续运营了 80 多年的文化产品”,你会发现,这场比较本质上是两个不同 IP 策略的长期结果对比:超人赢在定义了整个超级英雄题材…

2026/9/9 16:31:09

基于CNN的调制信号识别:MATLAB实现时频图分类实战

简介:本资源是一套面向通信工程与信号处理方向学习者、研究者的深度学习实践方案,聚焦调制信号自动检测与识别这一典型无线通信任务,解决传统方法依赖人工特征、低信噪比下性能下降等痛点。压缩包共12个文件(10.73MB)&…

2026/9/12 0:04:17

MATLAB仿生优化框架:长鼻浣熊算法多策略融合实现

简介:本资源是一份面向智能优化算法研究者与MATLAB初学者的仿生智能算法实践代码包,聚焦于长鼻浣熊优化算法(COA)的多策略改进与性能验证。针对传统COA易陷局部最优、收敛精度不足等问题,作者融合Circle映射初始化提升…

2026/9/12 0:04:17

【JAVA毕设源码分享】基于 JavaWeb 的校园一卡通管理系统的设计与实现 基于 JavaWeb 的校园卡业务管理系统(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/9/12 0:04:17

【JAVA毕设源码分享】基于 Java 的图书馆借阅管理平台的搭建与实现 基于 Java 的图书馆综合管理系统(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/9/12 6:29:36

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/10 15:19:50

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/12 6:37:43

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码