Data Engineering Zoomcamp 2025 第2周作业实战:用 Kestra 回填 2021 年 NYC 出租车数据

发布时间:2026/9/12 14:03:35
Data Engineering Zoomcamp 2025 第2周作业实战:用 Kestra 回填 2021 年 NYC 出租车数据 Data Engineering Zoomcamp 2025 第2周作业实战用 Kestra 回填 2021 年 NYC 出租车数据【免费下载链接】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 2025 学员第 2 周作业cohorts/2025/02-workflow-orchestration/homework.md展开以课程仓库中 Kestra 的 ETL Flow 源码为证据系统讲解如何把 2019/2020 年出租车数据管道扩展到 2021 年包含 Backfill 回填与手动执行两条路线、ForEach Subflow 循环挑战、6 道测验题的解析方法以及提交规范。读完本文你将能独立复现按年月拉取 NYC TLC 出租车 CSV → 上传 GCS → 写入 BigQuery的完整管道并掌握 Kestra 调度触发器、KV Store、变量渲染等核心概念在实战中的用法。作业背景与任务目标课程第 2 周使用 Kestra 构建 ETL 管道处理纽约市出租车与轿车委员会TLC的Yellow与Green两类出租车数据。课程视频中已演示了 2019 与 2020 两年的数据处理而本次作业的核心任务是将现有的 Flow 扩展到 2021 年数据。作业使用的数据集为green出租车数据集原始数据托管在 DataTalksClub 的nyc-tlc-data发布仓库中release tag 为green。为了得到可被wget直接下载的地址需要使用如下前缀原文档特别提示直接点击该链接会返回 404需配合文件名使用https://github.com/DataTalksClub/nyc-tlc-data/releases/download/green/文件名遵循{taxi}_tripdata_{year}-{month}.csv.gz的约定例如green_tripdata_2021-01.csv.gz。该约定与仓库中所有 Flow 的变量定义一一对应是理解整条管道的关键线索。上图展示的是作业所涉及的数据集年份与文件类型清单图片引用自课程 02-workflow-orchestration/images 目录与作业正文中的引用一致。前置条件搭建本周课程环境动手之前需要先把 Kestra 跑起来。本周模块的 README 给出了标准安装方式用 Docker Compose 同时启动 Kestra 服务器与配套的 Postgres 数据库Kestra 自身使用然后在浏览器访问http://localhost:8080打开 Kestra UI。如果你希望用命令行而非 UI 手动粘贴导入 FlowREADME 提供了基于 Kestra API 的批量导入方式形如curl -X POST http://localhost:8080/api/v1/flows/import -F fileUploadflows/06_gcp_taxi.yaml curl -X POST http://localhost:8080/api/v1/flows/import -F fileUploadflows/06_gcp_taxi_scheduled.yaml本周涉及的核心 Flow 全部位于 cohorts/2025/02-workflow-orchestration/flows 目录Flow 文件作用01_getting_started_data_pipeline.yaml入门演示HTTP 下载 → Python 转换 → DuckDB 查询02_postgres_taxi.yaml本地 Postgres 版出租车 ETL手动指定 year/month02_postgres_taxi_scheduled.yaml本地 Postgres 版调度 Flow含 Schedule 触发器03_postgres_dbt.yaml用 dbt 对 Postgres 中的原始数据做转换可选04_gcp_kv.yaml向 KV Store 写入 GCP 项目配置05_gcp_setup.yaml创建 GCS Bucket 与 BigQuery Dataset06_gcp_taxi.yamlGCS BigQuery 版出租车 ETL手动指定 year/month06_gcp_taxi_scheduled.yamlGCS BigQuery 版调度 Flow作业回填的推荐载体07_gcp_dbt.yaml用 dbt 对 BigQuery 数据做转换可选GCP 环境准备KV Store 与资源创建云上管道06_gcp_taxi.yaml等的 SQL 与上传逻辑大量引用 KV Store 中的键值。在运行任何 GCP Flow 之前需要先执行 04_gcp_kv.yaml把下列配置写入 KV Store注意将示例值替换为你自己的GCP_PROJECT_IDGCP 项目 ID示例kestra-sandboxGCP_LOCATION资源区域示例europe-west2GCP_BUCKET_NAMEGCS Bucket 名称必须是全局唯一GCP_DATASETBigQuery Dataset 名示例zoomcamp随后运行 05_gcp_setup.yaml它会通过io.kestra.plugin.gcp.gcs.CreateBucket创建 GCS BucketstorageClass: REGIONALifExists: SKIP并通过io.kestra.plugin.gcp.bigquery.CreateDataset创建 BigQuery Dataset。所有 GCP 插件共享的凭据与项目参数由 Flow 末尾的pluginDefaults统一注入pluginDefaults: - type: io.kestra.plugin.gcp values: serviceAccount: {{kv(GCP_CREDS)}} projectId: {{kv(GCP_PROJECT_ID)}} location: {{kv(GCP_LOCATION)}} bucket: {{kv(GCP_BUCKET_NAME)}}注意GCP_CREDS服务账号 JSON 属于敏感信息README 明确警告不要提交到 Git应像密码一样妥善保管。方案一利用 Backfill 功能回填 2021 年数据推荐作业的第一条提示非常明确复用调度 Flow 的 Backfill 功能。对应文件是 06_gcp_taxi_scheduled.yaml它与手动版06_gcp_taxi.yaml最大的区别在于时间不再来自用户输入而是来自触发器的日期。为什么调度 Flow 天然适合回填看它的变量定义06_gcp_taxi_scheduled.yamlvariables: file: {{inputs.taxi}}_tripdata_{{trigger.date | date(yyyy-MM)}}.csv gcs_file: gs://{{kv(GCP_BUCKET_NAME)}}/{{vars.file}} table: {{kv(GCP_DATASET)}}.{{inputs.taxi}}_tripdata_{{trigger.date | date(yyyy_MM)}} data: {{outputs.extract.outputFiles[inputs.taxi ~ _tripdata_ ~ (trigger.date | date(yyyy-MM)) ~ .csv]}}关键点是trigger.date在调度执行时它是当前触发时间而在Backfill 时它会变成你回填的那个历史日期。Kestra 的 Backfill 本质上是为日期范围内的每一天/每个月重新触发一次执行每次执行都带着对应的trigger.date于是文件名的年月部分yyyy-MM会被自动替换为回填的月份。这正是把管道扩展到历史数据零成本的原因。文件名、GCS 对象路径与 BigQuery 表名三者的年月格式并不完全一致值得注意变量渲染示例green / 2021-03说明vars.filegreen_tripdata_2021-03.csv带横杠的年月用于本地文件名与下载 URLvars.gcs_filegs://your-bucket/green_tripdata_2021-03.csvGCS 对象路径vars.tablezoomcamp.green_tripdata_2021_03下划线年月用于 BigQuery 临时表名Flow 末尾定义了两个 Schedule 触发器06_gcp_taxi_scheduled.yamltriggers: - id: green_schedule type: io.kestra.plugin.core.trigger.Schedule cron: 0 9 1 * * # 每月 1 日 09:00 UTC inputs: taxi: green - id: yellow_schedule type: io.kestra.plugin.core.trigger.Schedule cron: 0 10 1 * * # 每月 1 日 10:00 UTC inputs: taxi: yellow每个触发器通过inputs固化taxi类型因此回填时需要分别为yellow与green各执行一次。在 UI 中执行 Backfill 的步骤在 Kestra UI 打开 Flowzoomcamp.06_gcp_taxi_scheduled进入Triggers面板点击对应taxi触发器旁的Backfill按钮选择回填时间范围2021-01-01至2021-07-31——这是数据实际存在的区间2021 年 7 月之后 TLC 数据格式发生变化分别对yellow和green两个触发器各执行一次通过在触发器inputs中指定正确的taxi值建议勾选/添加标签backfill:true便于在 Execution 列表里区分回填产生的执行记录——这一点在 Flow 的description中已明确提示。回填期间每月会生成一次执行共 7 次 × 2 种车型 14 次执行每次执行会自动完成下载当月 CSV → 上传 GCS → 建外部表 → 写入月度临时表 → MERGE 进主表的完整链路。需要提醒的是本地 Postgres 版调度 Flow 02_postgres_taxi_scheduled.yaml 还额外声明了concurrency.limit: 1避免大数据量回填时并发执行相互竞争云上版本则可以放心地全量回填README 明确说明云环境存储与计算可无限扩展不会像本地机器那样耗尽资源。方案二手动逐月执行 ForEach / Subflow 循环挑战如果不想用 Backfill也可以使用手动版 Flow 06_gcp_taxi.yaml 逐月执行它在输入里显式声明了year与month两个SELECT输入且year输入特别开启了allowCustomValue: true源码中的注释写着 allows you to type 2021 from the UI for the homework因此在 UI 上可以直接输入2021再为 17 月各执行一次- id: year type: SELECT displayName: Select year values: [2019, 2020] defaults: 2019 allowCustomValue: true # allows you to type 2021 from the UI for the homework 手动路线下你需要为yellow和green各执行 7 次2021-01 至 2021-07共 14 次手动执行。挑战题用 ForEach Subflow 自动化循环作业给了一个进阶挑战找出如何在 ForEach 任务中循环 Year-Month 与 taxi 类型 的组合并用 Subflow 任务为每个组合触发一次 Flow 执行。思路是新建一个调度器Flow用io.kestra.plugin.core.flow.ForEach遍历组合列表在每次迭代中用io.kestra.plugin.core.flow.Subflow携带inputs调用现有的06_gcp_taxi或调度版Flow。下面是一个可供参考的示意结构注意这是基于 Kestra 官方任务类型编写的参考写法并非仓库现成文件你需要根据自己 Flow 的id与namespace调整id: 03_homework_foreach namespace: zoomcamp tasks: - id: for_each_combinations type: io.kestra.plugin.core.flow.ForEach values: - {taxi:yellow,year:2021,month:01} - {taxi:yellow,year:2021,month:02} # ... 其余月份与 green 组合 tasks: - id: run_taxi_flow type: io.kestra.plugin.core.flow.Subflow flowId: 06_gcp_taxi namespace: zoomcamp inputs: taxi: {{ taskrun.value | jq(.taxi) }} year: {{ taskrun.value | jq(.year) }} month: {{ taskrun.value | jq(.month) }}要点说明ForEach的values支持 JSON 数组每次迭代会把当前元素注入taskrun.valueSubflow通过flowIdnamespace指向被调用的子 Flow并通过inputs把参数透传下去被调用的子 Flow 执行成功后其结果如 GCS 对象、执行状态可在父 Flow 中通过{{ outputs.run_taxi_flow }}进一步引用。作业背后的管道原理从源码看完整 ETL 链路无论采用哪种方案最终落到 BigQuery 的写入逻辑都来自 Flow 本身。下面以 06_gcp_taxi_scheduled.yaml 为骨架拆解每个任务节点手动版 06_gcp_taxi.yaml 结构相同set_labelio.kestra.plugin.core.execution.Labels为执行打上file与taxi标签方便在 UI 中按文件/车型检索执行记录extractio.kestra.plugin.scripts.shell.Commands执行wget -qO- ...{{render(vars.file)}}.gz | gunzip {{render(vars.file)}}即流式下载 gzip 压缩的 CSV 并解压outputFiles声明*.csv为产出物供后续任务引用upload_to_gcsio.kestra.plugin.gcp.gcs.Upload把解压后的 CSV 上传到gs://{{kv(GCP_BUCKET_NAME)}}/{{vars.file}}if_yellow_taxi / if_green_taxiio.kestra.plugin.core.flow.If依据输入taxi分流各自执行 4 个 BigQuery 步骤建主表CREATE TABLE IF NOT EXISTS创建按日期分区PARTITION BY DATE(tpep_pickup_datetime)/DATE(lpep_pickup_datetime)的yellow_tripdata/green_tripdata主表建外部表CREATE OR REPLACE EXTERNAL TABLE直接指向 GCS 上的 CSVuris指向{{render(vars.gcs_file)}}并配置skip_leading_rows 1、ignore_unknown_values TRUE建月度临时表CREATE OR REPLACE TABLE ... AS SELECT从外部表读取用MD5(CONCAT(COALESCE(...)))基于关键字段VendorID、上下车时间、上下车地点生成unique_row_id并附上来源filenameMERGE 进主表MERGE INTO ... WHEN NOT MATCHED THEN INSERT按unique_row_id去重合并保证同一行程只入库一次purge_filesio.kestra.plugin.core.storage.PurgeCurrentExecutionFiles清理本次执行下载的中间文件避免撑爆 Kestra 存储。Green 与 Yellow 两套 SQL 的唯一差别在于字段集合Green 多出ehail_fee与trip_type且时间字段为lpep_pickup_datetime/lpep_dropoff_datetimeYellow 则使用tpep_pickup_datetime/tpep_dropoff_datetime且没有这两个字段。这正是作业中select the right service in thetaxiinput的原因——选错车型会导致外部表与目标表列数不匹配。本地 Postgres 版 02_postgres_taxi.yaml 的逻辑类似但数据装载使用io.kestra.plugin.jdbc.postgresql.CopyIn以header: true方式批量导入 staging 表再通过md5(COALESCE(...))生成unique_row_id后执行MERGE。注意 README 的提醒MERGE语句从 PostgreSQL 15 才引入因此本地练习必须使用 Postgres 15 或更高版本镜像否则会报语法错误。六道测验题解析作业包含 6 道选择题用于检验对工作流编排、Kestra 与 ETL 管道的理解。原文档题目与选项完整如下并附解题思路Q1执行 2020 年 12 月 Yellow 出租车数据时extract任务产出的yellow_tripdata_2020-12.csv未压缩文件大小是多少128.3 MiB / 134.5 MiB / 364.7 MiB / 692.6 MiB解题思路该文件大小来自extract任务的真实产出。可以在执行详情页的 Outputs 面板查看该 CSV 的大小或对下载解压后的文件执行ls -lh/du -h确认。这一题需要你实际运行后核对仓库源码中不存在该数值。Q2当输入taxigreen、year2020、month04时变量file渲染出的值是什么{{inputs.taxi}}_tripdata_{{inputs.year}}-{{inputs.month}}.csvgreen_tripdata_2020-04.csvgreen_tripdata_04_2020.csvgreen_tripdata_2020.csv解题思路本题可以直接从源码推出答案。在 06_gcp_taxi.yaml 中变量定义为file: {{inputs.taxi}}_tripdata_{{inputs.year}}-{{inputs.month}}.csv。{{inputs.taxi}}会先渲染为模板字符串本身Pebble 模板中{{ ... }}内的内容作为表达式求值将green、2020、04代入后得到green_tripdata_2020-04.csv即第二个选项。第一个选项是模板的原始定义而非渲染结果其余两个选项的年月顺序均与定义不符。Q32020 年全年所有 Yellow 出租车 CSV 文件共有多少行13,537,299 / 24,648,499 / 18,324,219 / 29,430,127Q42020 年全年所有 Green 出租车 CSV 文件共有多少行5,327,301 / 936,199 / 1,734,051 / 1,342,034Q52021 年 3 月 Yellow 出租车 CSV 文件有多少行1,428,092 / 706,911 / 1,925,152 / 2,561,031Q3Q5 解题思路这三题都需要用真实数据验证。推荐两条路径BigQuery 路径完成回填后对主表执行SELECT COUNT(*) FROM ... WHERE filename LIKE yellow_tripdata_2020-%按filename字段过滤该字段在临时表写入时已用{{render(vars.file)}}填充Q5 则过滤filename yellow_tripdata_2021-03.csv文件路径直接对解压后的 CSV 执行wc -l注意减去表头行。Q6如何为 Schedule 触发器把时区配置为纽约时间在 Schedule 触发器配置中新增timezone属性并设为EST在 Schedule 触发器配置中新增timezone属性并设为America/New_York在 Schedule 触发器配置中新增timezone属性并设为UTC-5在 Schedule 触发器配置中新增location属性并设为New_York解题思路Kestra 的 Schedule 触发器通过timezone属性指定时区且必须使用IANA 时区标识符如America/New_York而非EST、UTC-5这类缩写或偏移量写法也不存在location属性。例如配置cron: 0 9 1 * *timezone: America/New_York即表示纽约时间每月 1 日上午 9 点触发。仓库中两个调度 Flow 未显式声明timezone因此默认采用服务器/UTC 时区——这恰好是本题考察的知识点。注意Q1、Q3、Q4、Q5 的答案依赖实际运行结果本文不提供未经验证的数值Q2 可由 06_gcp_taxi.yaml 直接推导Q6 依据 Kestra Schedule 触发器的时区配置约定。提交要求与注意事项作业提交遵循以下要求来自原文档代码仓库提交表单末尾需要附上你的 GitHub 或其他公开代码托管仓库链接仓库中应包含解题代码如果方案不以文件形式存在例如纯 UI 操作请把步骤直接写进仓库的 README答案四舍五入如果没有完全匹配的选项选择最接近的一个提交方式与截止时间通过课程官方提交表单提交截止日期以表单页面为准原文档注明参考答案会在截止日期后补充因此完成作业时请以自己实际运行、验证得到的结果为准。常见问题排查来自模块 README结合本周模块 README 的排错建议做作业时常遇到三类问题Postgres 版本过低导致 MERGE 报错本地版 Flow 依赖 PostgreSQL 15 的MERGE语句请使用postgres:latest15 及以上镜像Kestra 镜像使用kestra/kestra:latest稳定版不要用develop分支镜像Linux 下连接host.docker.internal失败该域名在 Linux 上行为不同可改用 README 提供的一体化 Docker ComposeKestra 课程 Postgres pgAdmin 放同一文件在pluginDefaults中把地址改为容器名postgres_zoomcampBigQuery 报 CSV table references column position 17, but line contains only 14 columns这通常是 CSV 下载/上传不完整导致的列数不匹配并非 schema 问题。重新执行整个 Flow包括重新下载 CSV 并重新上传 GCS即可解决。总结本次作业的核心收获有三点其一理解Backfill 与trigger.date的配合机制——调度 Flow 因把年月写死在触发器日期中天然具备回填任意历史区间的能力只需在 UI 选择2021-01-01至2021-07-31并对yellow/green各执行一次其二掌握Kestra 的变量渲染{{vars.*}}、{{inputs.*}}、{{kv(...)}}、{{trigger.date | date(yyyy-MM)}}这是读懂与编写 Flow 的基础其三通过 ForEach Subflow 的挑战体会 Kestra 用声明式 YAML 编排循环触发子流程这类复杂调度场景的写法。配合本文对 06_gcp_taxi.yaml、06_gcp_taxi_scheduled.yaml 源码的逐节点拆解你可以在运行作业的同时把 Kestra 的 ETL 管道原理完整吃透。【免费下载链接】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),仅供参考

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询