ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Data Engineering Zoomcamp 实时流处理实战:用 PyFlink 从零搭建 Producer → Redpanda → Flink → PostgreSQL 出租车事件流水线

Data Engineering Zoomcamp 实时流处理实战:用 PyFlink 从零搭建 Producer → Redpanda → Flink → PostgreSQL 出租车事件流水线 Data Engineering Zoomcamp 实时流处理实战用 PyFlink 从零搭建 Producer → Redpanda → Flink → 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 2027 届第七模块「流处理Streaming」的开篇围绕仓库中 cohorts/2027/07-streaming/01-introduction.md 展开。该模块以一场基于纽约黄色出租车NYC yellow taxi行程数据的 PyFlink 流处理工作坊为核心从消息代理、生产者与消费者这些最基础的概念起步逐步加入数据库最终引入流处理框架 Flink端到端构建一条实时事件流水线Producer (Python) - Kafka (Redpanda) - Flink - PostgreSQL。读完本文你将掌握工作坊的整体架构、代码目录组织方式、前置环境准备以及贯穿全模块 14 个单元的学习路径并能在仓库中找到每一步对应的可运行源码。工作坊概览我们要构建什么这个工作坊的目标不是一次性交付成品而是从零逐步搭建一条实时流处理流水线让读者理解每个组件在流水线中的职责边界ProducerPython——把出租车行程事件发送到消息代理KafkaRedpanda——消息代理负责接收、存储并分发消息Flink——流处理框架对事件做透传、聚合等加工PostgreSQL——关系型数据库作为最终的持久化存储。整个链路可以用一行箭头概括来自 01-introduction.mdProducer (Python) - Kafka (Redpanda) - Flink - PostgreSQL它的价值在于把一条「实时事件」从生产端一路送到「持久化关系存储」的完整生命路径走通。前几步broker、producer、consumer解决的是事件如何流动加入数据库后解决的是事件如何落盘而引入 Flink 解决的则是事件如何在流动过程中被可靠地计算——这正是后续单元层层递进的主线。前置条件与环境准备工作坊对环境的依赖极简只需要三类工具前置条件用途备注Docker 与 Docker Compose运行 Redpanda、Flinkjobmanager/taskmanager和 PostgreSQL 服务所有基础设施均以容器方式启动uvPython 项目管理与依赖安装仓库使用pyproject.tomluv.lock管理 Python 依赖uv sync即可还原环境SQL 客户端连接 PostgreSQL 验证落库结果可用uvx pgcli无需单独安装、DBeaver、pgAdmin 或 DataGrip其中 SQL 客户端的最小可用方案是uvx pgcli——uvx是 uv 自带的临时运行工具无需先安装 pgcli 即可直接调用非常适合快速验证。仓库代码结构参考实现与现场演示代码工作坊代码全部位于 cohorts/2027/07-streaming/code/其中包含两套互补的代码code/参考代码工作坊的最终完整实现按src/组织成清晰的模块化结构code/src/producers/ —— 生产端脚本code/src/consumers/ —— 消费端脚本控制台输出与 PostgreSQL 落库两个版本code/src/job/ —— Flink 作业透传作业与聚合作业code/src/models.py —— 共享数据模型code/live/工作坊现场代码讲师在直播工作坊中逐步敲出的版本结构与参考代码对应同样包含src/job/、src/producers/和notebooks/。这种参考实现 现场演示的双目录设计意味着你可以跟随后续单元从零手写每一步也可以直接研读既有文件、按命令运行验证。工作坊的进阶路线是边看边写两者对照能最快理解每个文件为何长成现在的样子。学习路径后续 14 个单元如何串起整条流水线01-introduction.md明确指出接下来的单元会从零构建一切。完整路径如下见 README.mdPyFlink: Stream Processing Workshop —— 本指南总览Redpanda - a Kafka-compatible broker —— 启动消息代理Produce messages to Kafka —— 用 Python 生产消息Consume messages with Python —— 用 Python 消费消息Save events to PostgreSQL —— 事件落库Why Flink? —— 为什么需要流处理框架The Flink image and services —— Flink 镜像与作业管理/任务管理服务The pass-through Flink job —— 第一个 Flink 作业透传Offsets - earliest vs latest —— 偏移量语义Aggregation with tumbling windows —— 翻滚窗口聚合Late events and upserts —— 迟到事件与更新插入Understanding window types —— 窗口类型深入Cleanup —— 资源清理QA —— 答疑这套路径本身就是一个递进的教学设计第 25 单元用手写代码解决消息流动 落库到第 5 单元结尾抛出如果要按时间窗口聚合、要容错、要并行、要写多个目标怎么办的痛点自然引出第 6 单元之后的 Flink 主题。仓库还提供了配套的 homework.md 作为课后练习。深入源码四个关键组件如何协作下面结合 code/ 下的真实实现逐层拆解流水线各环节。1. 共享数据模型src/models.py生产端与消费端共用同一个Ridedataclass见 code/src/models.py它定义了事件的最小 schemadataclass class Ride: PULocationID: int DOLocationID: int trip_distance: float total_amount: float tpep_pickup_datetime: int # epoch milliseconds两个细节值得注意tpep_pickup_datetime用 epoch 毫秒整数表示而不是字符串时间戳。这是因为后续 Flink 端要把它转换TO_TIMESTAMP_LTZ成 TIMESTAMP 并作为事件时间使用整数时间戳是跨语言、跨系统最无歧义的交换格式配套的ride_from_row()负责把 pandas DataFrame 的一行转成Rideint(row[tpep_pickup_datetime].timestamp() * 1000)而ride_deserializer()则把 Kafka 交付的 JSON 字节串还原为Ride对象构成序列化/反序列化的闭环。2. 生产端src/producers/producer_realtime.py参考代码中除了基于真实 parquet 数据的producer.py还有一个专为工作坊设计的实时生产者 code/src/producers/producer_realtime.py。它模拟真实出租车流量的关键特性基于真实热点的上车/下车地点池PICKUP_LOCATIONS取自纽约 TLC 出租车分区 ID1263如 East Village(79)、JFK 机场(132)、时代广场(230) 等让模拟数据贴近真实分布~20% 概率生成迟到事件晚 310 秒make_ride(delay_secondsdelay)会把时间戳回拨这是后续迟到事件单元11-late-events-and-upserts.md专门要处理的场景通过KafkaProducer(bootstrap_servers[localhost:9092], value_serializerride_serializer)写入ridestopic每 0.5 秒发一条CtrlC时flush()并打印已发送总数。3. 消费端src/consumers/仓库提供两个消费端体现先看数据再存数据的教学节奏code/src/consumers/consumer.py把事件打印到控制台。关键参数是auto_offset_resetearliest从 topic 起点开始读与group_idrides-console消费组各自独立跟踪偏移量code/src/consumers/consumer_postgres.py用psycopg2把事件 INSERT 进 PostgreSQL 的processed_events表conn.autocommit True让每条 INSERT 即时提交无需手动commit()。它使用独立的消费组rides-to-postgres因此与终端消费者各自都能读到全量消息。第 5 单元在完成手写落库后会用一段还缺什么的反思清单窗口聚合逻辑、崩溃后的 offset 管理、多实例并行、多 sink 连接器点明手写方案的边界——这正是 Flink 登场的理由。4. 基础设施编排docker-compose.yml工作坊的全部服务由一份 code/docker-compose.yml 编排包含四个服务服务镜像/构建关键配置对外端口redpandaredpandadata/redpanda:v25.3.9单节点、双 Kafka 监听器9092外部 Kafka 协议、29092容器内 Kafka 协议、8082/28082HTTP Proxyjobmanager./Dockerfile.flink构建为pyflink-workshopRPC 端口 6123、Web UI 8081、JVM 1600m8081Flink Web UItaskmanagerpyflink-workshop复用镜像15 个任务槽、默认并行度 3、JVM 1728m6121/6122postgrespostgres:18用户/密码/库均为postgres5432关于 Redpanda 的双监听器设计详见 02-redpanda.mdKafka 客户端采用两步连接——先连 bootstrap 服务器获取集群元数据再按 broker 返回的advertised 地址建立实际数据连接。因此容器内连接如 Flink 访问redpanda:29092与宿主机连接Python 访问localhost:9092必须分别广播地址否则必然有一方连不上。这也是后续 Flink 作业中properties.bootstrap.servers redpanda:29092的来源。5. Flink 镜像与配置Dockerfile.flink与flink-config.yamlFlink 运行环境不是默认镜像而是工作坊定制构建的code/Dockerfile.flink 基于flink:2.2.0-scala_2.12-java17通过 uv 安装 Python 3.12 使 Flink 具备 PyFlink 能力并预下载了四个关键连接器 jarflink-json、flink-sql-connector-kafka、flink-connector-jdbc-core、flink-connector-jdbc-postgres以及postgresqlJDBC 驱动——这正是 Flink 作业能读写 Kafka/PostgreSQL 的底层依赖code/flink-config.yaml 相对默认配置做了两处关键修改为 PyFlink 增加了taskmanager.memory.jvm-metaspace.size: 512mPython 作业需要更大元空间并移除了 JRE 中不存在的jdk.compiler模块--add-exports项以避免每次命令告警。6. 两个 Flink 作业透传与聚合src/job/下的两个作业构成了 Flink 部分的进阶主线透传作业code/src/job/pass_through_job.py以connector kafka声明ridestopic 为源表scan.startup.mode latest-offset即只读新消息以connector jdbc声明 PostgreSQL 的processed_events为汇表用一句INSERT INTO ... SELECT完成 Kafka → Postgres 的搬运动作期间用TO_TIMESTAMP_LTZ(tpep_pickup_datetime, 3)把 epoch 毫秒转成 TIMESTAMP作业每 10 秒启用一次 checkpointenv.enable_checkpointing(10 * 1000)保证故障恢复聚合作业code/src/job/aggregation_job.py在源表上声明事件时间列与水印WATERMARK for event_timestamp as event_timestamp - INTERVAL 5 SECOND容忍 5 秒迟到以TUMBLE(TABLE events, DESCRIPTOR(event_timestamp), INTERVAL 1 HOUR)做 1 小时翻滚窗口GROUP BY window_start, PULocationID统计COUNT(*)趟数与SUM(total_amount)营收汇表以(window_start, PULocationID)为 PRIMARY KEY 支持 upsert并将scan.startup.mode改为earliest-offset以重放全量数据、并行度设为 3。对照两个作业的差异就能直观理解本模块后半部分的核心概念偏移量语义latest vs earliest、事件时间与水印、翻滚窗口聚合、迟到事件与 upsert。小结与后续学习建议本模块把流处理知识组织成一条先动手、再反思、后上框架的学习曲线先用 Python Redpanda PostgreSQL 手搓一条可用但脆弱的事件链路体会其痛点再引入 Flink 用声明式 SQL 优雅地解决窗口聚合、容错与并行问题。仓库中的 code/ 是每一步的最终形态homework.md 可用于检验掌握程度而 theory/ 与 extras/ 则提供 Kafka 原理与历年 Python/PyFlink 示例作为课外延伸。建议按单元顺序边写边跑启动docker compose up redpanda postgres -d运行生产者与消费者验证消息流动再构建 Flink 镜像依次提交透传与聚合作业最后用uvx pgcli查询processed_events与processed_events_aggregated两张表直观看到实时计算的结果。【免费下载链接】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),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进