ARTICLE · INTELLIGENCE

战地情报 · 详情页

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

Prefect 工作流编排实战:Data Engineering Zoomcamp 第 2 周作业全解析(从 Web 加载、Cron 调度到 GitHub 存储与通知)

Prefect 工作流编排实战:Data Engineering Zoomcamp 第 2 周作业全解析(从 Web 加载、Cron 调度到 GitHub 存储与通知) Prefect 工作流编排实战Data Engineering Zoomcamp 第 2 周作业全解析从 Web 加载、Cron 调度到 GitHub 存储与通知【免费下载链接】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 2023 届第 2 周Workflow Orchestration的课后作业为骨架逐题拆解如何用Prefect编排一条完整的出租车数据管道从 Web 下载 CSV 加载到 Google Cloud StorageGCS、用 Cron 调度 flow、把 GCS 中的数据加载进 BigQuery、把 flow 代码存进 GitHub 仓库供团队协作再到配置 Email/Slack 通知与 Secret 块。读完你不仅能独立完成这套作业还能掌握 flow、task、deployment、storage block、automation、secret block 等 Prefect 核心概念在一套真实数据管道中的落地方式。作业背景工作流编排与可观测性训练2023 届第 2 周的主题是工作流编排Workflow Orchestration官方配套的视频课程与代码仓库围绕 Prefect 展开覆盖以下核心能力详见 cohorts/2023/week_2_workflow_orchestration/README.mdPrefect 基础概念flow、task、blocks、collections、Orion UIETL 与 GCPFlow 1 将数据写入 GCSFlow 2 从 GCS 加载到 BigQuery参数化 flow 与 Pydantic 参数校验、本地创建 deployment、启动 Prefect Agent、运行 flow调度Cron、flow 代码存储本地 / Docker / GitHub、在 Docker 中运行 taskPrefect Cloud 工作区与通知Automation / Notification。本次作业的目标非常明确——familiarise users with workflow orchestration and observation让学习者熟悉工作流编排与可观测性不仅要让管道能跑还要能定时跑、能看日志、能被通知、能安全地管理密钥。在动手前先准备好环境本地安装 Prefect含prefect-gcp等 collection拥有一个 GCP 项目并开启 GCS 与 BigQuery 服务完成gcloud auth application-default login或配置服务账号密钥。数据源是 DataTalksClub 发布的 NYC TLC 出租车数据集按{service}_tripdata_{year}-{month}.csv.gz命名的压缩 CSV服务包括 green / yellow / fhv。前置理解参考仓库中的 web_to_gcs 脚本作业反复引用视频课中的etl_web_to_gcs.py脚本作为起点。当前仓库虽然没有 2023 届的 Prefect 版脚本但保留了功能等价的非编排版参考实现03-data-warehouse/extras/web_to_gcs.py 与带进度条的增强版 web_to_gcs_with_progress_bar.py。阅读这两个脚本能帮你快速理解“Web → 本地 → GCS”这一经典 E 阶段应包含哪些步骤拼接并下载 CSVinit_url https://github.com/DataTalksClub/nyc-tlc-data/releases/download/按月循环构造{service}_tripdata_{year}-{month}.csv.gz并用requests下载增强版通过gzip.open先统计总行数、按 1 MB 分块流式下载并显示进度条转成 Parquet 并固定数据类型脚本通过dtypes字典把VendorID、PULocationID、DOLocationID、payment_type等指定为Int64把store_and_fwd_flag指定为string把金额类字段指定为float64再按 service 区分tpep_pickup_datetimeyellow与lpep_pickup_datetimegreen等日期列保证 Parquet 列类型直接可用与第 1 周ingest_data.py的做法一致上传到 GCSupload_to_gcs使用google.cloud.storage的blob.upload_from_filename上传对象路径形如green/green_tripdata_2020-01.parquet。把这三步分别包进task再在外面套一个flow就是作业要求的etl_web_to_gcs.py的基本形态。运行方式可参考 03-data-warehouse/extras/README.mduv sync安装依赖后执行uv run python web_to_gcs.py。问题 1把 2020 年 1 月 Green 出租车数据加载到 GCS题目以etl_web_to_gcs.py为模板新建一个 flow把2020 年 1 月的 Green 出租车 CSV 数据集加载进 GCS 并运行然后查看日志确认该数据集有多少行。选项447,770 / 766,792 / 299,234 / 822,132。解题思路这一步完全不涉及调度和部署纯粹是“让第一个 flow 跑通”。参考上面仓库中的web_to_gcs逻辑只需把目标月份固定为2020-01、service 固定为greenflow 内用requests下载green_tripdata_2020-01.csv.gz用 pandas 按固定dtypes与parse_dates[lpep_pickup_datetime, lpep_dropoff_datetime]读取写入 Parquet上传到gs://{你的bucket}/green/green_tripdata_2020-01.parquet在 flow 中print出len(df)行数并给flow装饰器加log_printsTrue这样打印内容才会进入 Prefect 的运行日志。验证方式在 Prefect Orion UI 中打开该 flow run查看日志输出或在终端观察 print即可读到该数据集的行数。答案447,770 行依据该课程官方解答也可通过运行你自己的 flow 查看日志交叉验证。问题 2用 Cron 表达式做月度调度题目利用etl_web_to_gcs.py中的 flow创建一个每月 1 日凌晨 5 点UTC运行的 deployment。正确的 cron 表达式是哪个选项0 5 1 * */0 0 5 1 */5 * 1 0 */* * 5 1 0。背景知识Cron 五段式语法。标准 cron 表达式由 5 个字段组成顺序为字段分钟小时日月星期取值范围0–590–231–311–120–70 和 7 都表示周日把需求“每月 1 日、凌晨 5 点UTC”逐字段翻译分钟 0小时 5日 1月 *任意月星期 *任意星期几得到0 5 1 * *。在 Prefect 中创建带调度的 deployment可用 CLI 的--cron参数prefect deployment build etl_web_to_gcs.py:etl_web_to_gcs -n monthly-first-day --cron 0 5 1 * * --storage-block local-file-system/dev也可以直接向Deployment.build_from_flow/serve传入schedules[CronSchedule(cron0 5 1 * *, timezoneUTC)]。注意明确指定timezoneUTC否则会按服务器本地时区解释“5 点”。答案0 5 1 * *其余选项的语义分别是0 0 5 1 *会把“5”放到“日”字段变成每月 5 日零点5 * 1 0 *表示每周一星期 0每小时的第 5 分钟运行* * 5 1 0每秒钟都会命中不是合法调度。问题 3从 GCS 加载到 BigQuery——参数化 flow 与本地子进程部署题目以etl_gcs_to_bq.py为起点改造脚本实现从 GCS 抽取Extract并加载Load进 BigQuery不填充也不删除缺失值即只做 ETL 中的 E 和 L不做 Transform。要求主 flow 打印脚本处理的总行数并让 flow 装饰器开启日志输出将入口 flow参数化接受months列表、year、taxi_color创建一个本地子进程 本地 flow 代码存储默认值的 deployment提前把 Yellow 出租车2019 年 2 月和2019 年 3 月的 Parquet 文件放进 GCS运行 deployment 把这两月数据追加进 BigQuery 表。解题要点只做 E 与 L读取 Parquet 时不要做dropna()或fillna()直接df.to_gbq(...)或经prefect-gcp的BigQueryWarehouseblock 写入保持缺失值原样入库参数化flow 签名形如flow(log_printsTrue, nameetl_gcs_to_bq)、def etl_gcs_to_bq(months: list[str], year: int, taxi_color: str)在函数体内按(year, month)拼接 GCS 对象路径gs://{bucket}/{color}/{color}_tripdata_{year}-{month}.parquet逐月读取、累加行数并print总行数。参数校验由 Prefect 内置的 Pydantic 完成list[str]、int这类类型注解会被自动强校验本地子进程部署默认值不指定任何 storage 和 infrastructure 时Prefect 默认 flow 代码存本地、运行在本地子进程中。构建命令为prefect deployment build etl_gcs_to_bq.py:etl_gcs_to_bq -n gcs-to-bq-append -q default prefect agent start -q default随后在 UI 中点击 Run或在 CLI 用prefect deployment run触发并传入参数{months: [02, 03], year: 2019, taxi_color: yellow}数据准备用第 1 题的脚本先把 yellow 2019-02、2019-03 的 Parquet 上传到 GCS仓库脚本 web_to_gcs.py 注释中web_to_gcs(2019, yellow)即是此类调用。验证方式观察 flow run 日志中打印的总处理行数。答案14,851,920 行即 Yellow 2019 年 2 月约 6,914,071 行与 3 月约 7,937,849 行的合计依据课程官方解答跑完自己的 deployment 后日志数字应与之吻合。问题 4用 GitHub Storage Block 托管 flow 代码题目团队协作时希望 flow 代码存在 GitHub 仓库里让 Prefect 从仓库读取代码而非本地存储或打进 Docker 镜像。创建一个GitHub storage blockUI 或 Python 代码均可并用它创建 deployment在本地子进程中运行处理2020 年 11 月的 Green 数据。注意代码要自己 push 到 GitHubPrefect 不会替你推送。解题思路存储块Storage Block是什么Prefect 的 block 是“配置 凭据”的可复用封装。github类型的 block 保存仓库 URL、分支、访问令牌私有仓库需要github-token运行时 Prefect 会按 block 中记录的路径拉取 flow 代码创建方式二选一UI 方式导航到 Blocks →→ 选择 GitHub填入仓库地址、分支如main、flow 文件路径保存并复制生成的块名如github/my-team-repo。代码方式from prefect_github import GitHubCredentials, GitHubRepository credentials GitHubCredentials(tokenghp_your_token) block GitHubRepository( namemy-team-repo, repositoryhttps://github.com/your-org/prefect-zoomcamp, referencemain, credentialscredentials, ) block.save(overwriteTrue)把 deployment 指向该存储块prefect deployment build flows/04_homework.py:flow_gcs_to_bq -n gh-stored --storage-block github/my-team-repo prefect deployment apply flow-deployment.yamldeployment 构建时不需要本地代码路径运行前 Prefect 会从 GitHub 拉取运行不指定 infrastructure 即默认本地子进程启动 agent 后手动触发 deployment参数为{months: [11], year: 2020, taxi_color: green}。答案88,605 行Green 2020 年 11 月数据依据课程官方解答。问题 5Email 或 Slack 通知——让失败/成功可感知题目数据管道出问题时最好能被及时通知。二选一方案 APrefect Cloud Email注册免费的 Prefect Cloud 账号app.prefect.cloud按 UI 指引连接工作区创建一条Automation当 flow run 进入 Completed 状态时给自己发邮件。然后运行第 4 题的 deployment参数换成Green 2019 年 4 月的数据检查邮箱是否收到通知方案 BSlack Webhook用 Prefect Cloud 的 Automation 或自托管 Orion 的Notification在 flow run 进入 Completed 状态时向 Slack 工作区发送消息。Webhook URL 既可以使用课程提供的临时邀请链接所对应的工作区也可以从你自己创建的 Slack 工作区 Slack App 中获取更推荐因为课程提供的临时链接有时效且有人数上限。关键概念区分Prefect Cloud 中这类基于事件的响应叫Automation支持“发生/未发生某事件”等复杂触发条件而自托管 Orion 服务端对应的功能叫Notification。两者在配置界面上都需要指定触发事件如PrefectFlowRun.Completed/PrefectFlowRun.Failed、通知渠道邮件或 incoming webhook与目标地址。验证方式完成配置后运行 flow等待运行结束确认收到 Email 或 Slack 消息。答案125,268 行Green 2019 年 4 月数据依据课程官方解答。问题 6Prefect Secret 块——敏感信息的加密存储与 UI 遮蔽题目Prefect 的Secret 块提供数据库中加密存储与 UI 中的遮蔽obfuscation能力。在 UI 中创建一个存有10 位假密码的 Secret 块用于模拟连接第三方服务创建完成后的下一页 UI 上密码会显示为多少个星号*选项5 / 6 / 8 / 10。解题思路与原理Secret 块的值在数据库中以加密形式保存写入后便不可回读明文。创建时输入 10 位密码如1234567890保存后回到块详情页Prefect 只显示定长的遮蔽串且遮蔽长度与明文长度无关这是为了同时隐藏“值的长度”这一信息。因此你看到的星号数量既不是 10也不是按密码位数 1:1 显示的。答案8 个星号Prefect UI 对 Secret 值统一渲染为 8 个*依据课程官方解答与 Prefect UI 的遮蔽行为可在自己的 Prefect 实例中创建后目测验证。提交作业与常见踩坑提交方式填写官方作业提交表单可多次提交以最后一次提交为准2023 届的截止时间是 2 月 8 日周三22:00 CET常见踩坑提示Q3 忘记在flow上加log_printsTrue导致行数没有进入日志Q4 忘记把最新代码 push 到 GitHubPrefect 拉到的是旧版本Q4 私有仓库未配置 GitHub tokenagent 拉取代码时 401Q5 使用临时 Slack 邀请链接前未注意其 90 天有效期与 400 人上限建议自建 Webhook调度时未显式指定timezone导致“凌晨 5 点”被按本地时区解释。小结通过这 6 道题你实际上走完了一条生产级数据管道从“单次执行”到“可调度、可协作、可通知、可加密管理密钥”的完整进化路径web_to_gcsE→ cron 调度 → gcs_to_bqL→ GitHub 代码托管 → 事件通知 → Secret 管理。这恰恰对应了本课程对工作流编排工具的核心诉求——不仅仅是“能跑通”更是要“跑得可靠、看得清楚、失败可知、凭据安全”。如果想温习本作业依托的课程大纲与视频章节可回到 cohorts/2023/week_2_workflow_orchestration/README.md想复用不带编排器的 Web→GCS 参考实现可查看 03-data-warehouse/extras/web_to_gcs.py 及其 README。【免费下载链接】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

延伸阅读

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