Data Engineering Zoomcamp 实战用 Kestra Schedule 触发器调度出租车数据管道并回填历史数据本地 Postgres 环境【免费下载链接】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 课程「Local DB: Learn Scheduling and Backfills」章节展开讲解如何为上一章的 NYC 出租车Yellow / Green Taxi数据加载管道加上**定时调度Scheduling与历史回填Backfill**能力。读完本文你将掌握 Kestra Schedule 触发器的 cron 配置与trigger.date动态变量机制学会在 Kestra UI 中安全地回填 2019 年 Green Taxi 的 12 个月历史数据并理解整个管道为何能做到按时间片重跑而不产生重复数据。为什么需要调度与回填在 07-load-taxi-data-to-postgres 中我们已经通过手动传入taxi/year/month输入项把某个具体月份的出租车 CSV 数据加载进了本地 Postgres。但生产环境中的数据管道不能永远靠人肉点击触发——我们需要自动调度让管道按照固定时间规律自动运行无需人工干预回填Backfill当管道或表结构上线时用同一套代码把过去的历史数据补齐保证时间序列数据的完整性。本节的做法是在上一章 04_postgres_taxi.yaml 的基础上去掉year/month两个手工输入项改为从触发器的执行时间自动推导年月再挂上两个 Schedule 触发器。完整代码见 05_postgres_taxi_scheduled.yaml。考虑到本地磁盘与数据库资源有限回填示例只处理 Green Taxi 在 2019 年全年12 个月的数据而非全量 Yellow Green 历史数据。调度版流程的整体结构调度版流程05_postgres_taxi_scheduled.yaml与手动版相比核心差异集中在三处对比维度手动版04_postgres_taxi.yaml调度版05_postgres_taxi_scheduled.yaml时间输入手工选择year、month由trigger.date自动推导年月触发方式UI 手动 Execute两个 Schedule 触发器自动触发并发控制无concurrency.limit: 1防止重叠执行数据表/文件命名基于输入的年月基于触发器日期的年月表结构与加载逻辑相同相同输入与变量时间片从哪来inputs: - id: taxi type: SELECT displayName: Select taxi type values: [yellow, green] defaults: yellow variables: file: {{inputs.taxi}}_tripdata_{{trigger.date | date(yyyy-MM)}}.csv staging_table: public.{{inputs.taxi}}_tripdata_staging table: public.{{inputs.taxi}}_tripdata data: {{outputs.extract.outputFiles[inputs.taxi ~ _tripdata_ ~ (trigger.date | date(yyyy-MM)) ~ .csv]}}关键点只保留taxi一个输入项yellow为默认值trigger.date是 Kestra 在触发器含调度触发与回填触发的执行中自动注入的变量代表本次执行对应的时间点Pebble 模板过滤器date(yyyy-MM)把它格式化为2019-01这样的年月字符串进而拼出文件名green_tripdata_2019-01.csvdata变量通过outputs.extract.outputFiles[...]引用 extract 任务产生的 CSV 文件作为后续CopyIn的数据源。并发限制回填安全的保障concurrency: limit: 1concurrency.limit: 1保证同一时刻最多只有一个执行在跑。回填 12 个月意味着连续触发 12 次执行如果并发失控、多个执行同时去TRUNCATE同一个 staging 表数据就会被互相破坏。这个限制在回填场景下是必要的保护措施。任务流与手动版完全一致的加载链路调度版复用了手动版的全部任务逻辑通过if_yellow_taxi/if_green_taxi两个条件分支按车型分流每个分支依次执行extract用wget -qO- ... | gunzip流式下载并解压对应月份的 CSV*_create_table/*_create_staging_table创建目标表与 staging 表IF NOT EXISTS*_truncate_staging_table清空 staging 表保证当前批次数据干净*_copy_in_to_staging_table用 PostgreSQL 的CopyIn批量灌入 CSVheader: true跳过表头*_add_unique_id_and_filename用md5()对VendorID pickup dropoff PULocationID DOLocationID fare_amount trip_distance拼出的指纹生成unique_row_id并记录来源filename*_merge_dataMERGE INTO把 staging 数据并入最终表ON T.unique_row_id S.unique_row_id只插入不存在的行purge_files清理本次执行产生的临时文件。其中unique_row_idMERGE的组合让管道天然幂等——同一时间片无论被调度触发多少次、被回填重复跑多少次都不会产生重复记录。这正是调度 回填能够安全落地的基石。数据库连接信息通过pluginDefaults统一注入无需在每个任务里重复书写pluginDefaults: - type: io.kestra.plugin.jdbc.postgresql values: url: jdbc:postgresql://pgdatabase:5432/ny_taxi username: root password: rootSchedule 触发器cron 表达式与输入注入调度能力的载体是流程末尾的triggers段triggers: - id: green_schedule type: io.kestra.plugin.core.trigger.Schedule cron: 0 9 1 * * inputs: taxi: green - id: yellow_schedule type: io.kestra.plugin.core.trigger.Schedule cron: 0 10 1 * * inputs: taxi: yellow逐字段解读type: io.kestra.plugin.core.trigger.ScheduleKestra 核心调度触发器按 cron 表达式产生定时执行cron: 0 9 1 * *标准五段 cron含义为每月 1 日 09:00 UTC分0、时9、日1、月、周0 10 1 * *则是每月 1 日 10:00 UTCinputs: {taxi: green}/{taxi: yellow}为每次触发注入taxi输入项实现一个流程、两个调度实例——Green 和 Yellow 各自按自己的时间表运行。需要留意的是本节课程文档文字描述为每天 9 AM UTC 调度但仓库中该 YAML 实际采用的 cron 是0 9 1 * *每月 1 日 09:00 UTC。从数据结构看这更贴合实际TLC 数据按年月切片成{taxi}_tripdata_{yyyy-MM}.csv每月 1 日跑一次正好加载上一个自然月的新文件配合 MERGE 幂等逻辑也不会产生重复。如果你确实希望改成每天运行只需把 cron 换成0 9 * * *即可管道对重复执行是安全的。相关概念触发器、执行、变量、输入可回看 04-kestra-concepts。调度执行的时间线每次调度触发时Kestra 会为执行生成trigger.date默认为触发时刻UTC 时区用inputs注入taxi渲染所有{{render(vars.*)}}变量开始执行任务流。因此每次调度执行天然对应一个确定的年月时间片互不干扰。回填实战在 Kestra UI 中补齐 2019 年 Green Taxi 数据回填的本质是把调度触发器拨回到过去的时间点逐片重放执行。Kestra 在 UI 中直接支持这一操作GCP 章节文档也明确写着 You can backfill historical data directly from the Kestra UI。为什么只回填 Green 2019本地开发环境跑在单机 Docker 上Postgres 与磁盘空间都有限。Yellow 数据量远大于 Green若把两个车型多年历史全部回填本地资源可能吃不消。因此本节的实践范围收敛为Green Taxi、2019 年 1–12 月正好 12 个时间片。操作步骤确保已把流程 05_postgres_taxi_scheduled.yaml 导入 Kestra可从 UI 编辑器粘贴或用curl -X POST -u adminkestra.io:Admin1234! http://localhost:8080/api/v1/flows/import -F fileUploadflows/05_postgres_taxi_scheduled.yaml经 API 导入见 03-installing-kestra打开该流程详情页进入Triggers标签页在green_schedule触发器上点击Backfill设置回填时间范围开始2019-01-01、结束2019-12-31Kestra 会按触发器的 cron 粒度切分即每月 1 日一个时间片共 12 次执行确认时区默认 UTC与trigger.date的处理保持一致在输入覆盖中确认taxi: green提交回填Kestra 会依次生成 12 个执行每个执行对应green_tripdata_2019-0X.csv。用标签追踪回填执行流程的 description 字段给出一条实践建议在 UI 为回填产生的执行添加backfill:true标签。数据管道在生产中往往同时存在定时增量执行与历史回填执行给回填执行打上标签后在 Executions 列表按标签过滤、或在后续血缘分析中区分两类执行都会非常方便。回填执行与普通调度执行走的是同一条代码路径trigger.date被设为回填的日期因此每个执行会精确加载对应月份的文件并通过 staging MERGE 幂等合并。如何验证调度与回填结果Executions 页面回填提交后应看到 12 个执行逐个对应 2019 年 1 月至 12 月concurrency.limit: 1保证它们串行执行、互不踩踏Gantt 视图点开任意执行查看 Gantt 标签可直观看到 extract → create table → truncate → copy in → merge → purge 各任务的耗时与依赖关系06-getting-started-pipeline 中也建议用 Gantt 与 Logs 检查执行情况Logs 标签确认 extract 任务下载的是正确的green_tripdata_2019-XX.csvCopyIn成功且无报错Postgres 侧校验用 pgAdmin仓库 docker-compose 中运行在 8085 端口连接ny_taxi库查询public.green_tripdata按filename分组应能看到 12 个文件对应的数据且SELECT COUNT(*) FROM public.green_tripdata与各月 CSV 行数之和一致无重复因为 MERGE 只插新行重复运行安全性对同一时间片再次手动触发执行行数应保持不变——这直接验证了管道的幂等性。从本地到云端下一站本节在本地 Postgres 上完成了调度 回填闭环后续章节会把这套模式原样迁移到云上10-setup-google-cloud-platform 完成 GCP 环境与密钥KV/Secret配置11-load-taxi-data-to-bigquery 将数据加载目标切换到 GCS BigQuery12-schedule-and-backfill-full-dataset 在云端回填全部Yellow 与 Green 历史数据——其流程 09_gcp_taxi_scheduled.yaml 使用完全相同的 cron 与回填机制只是加载目标换成了可无限扩展的云存储与云数仓不再受本地磁盘限制。小结本节的本质是用 Kestra 的Schedule 触发器 trigger.date动态变量 幂等合并三个机制把一个手动数据管道升级为可调度、可回填、可重放的生产级管道。无论本地还是云端这套模式完全一致调度负责新数据自动入库回填负责历史数据按时间片补齐而 staging 表 unique_row_idMERGE保证无论跑多少遍数据都不会重复。【免费下载链接】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),仅供参考