ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

用 Bruin 与 DuckDB 构建端到端 NYC Taxi 数据管道:三层架构实战(Data Engineering Zoomcamp Module 5.3)

用 Bruin 与 DuckDB 构建端到端 NYC Taxi 数据管道:三层架构实战(Data Engineering Zoomcamp Module 5.3) 用 Bruin 与 DuckDB 构建端到端 NYC Taxi 数据管道三层架构实战Data Engineering Zoomcamp Module 5.3【免费下载链接】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 第 5 模块的第 3 课05-data-platforms/notes/03-nyc-taxi-pipeline.md为核心完整演示如何用数据平台 Bruin 以本地 DuckDB 为数据库从零搭建一条摄取ingestion→ 清洗staging→ 报表reports三层的 NYC Taxi 端到端管道。你将掌握bruin init初始化项目、.bruin.yml与pipeline.yml的配置、Python / SQL / Seed 三类 asset 的编写、time_interval等物化策略、数据质量检查、变量注入与bruin validate/run/query命令行操作最终能在本机跑通一条具备依赖编排、增量处理和血缘可视化的完整数据管道。整体架构三层的 ELT 管道Bruin 把数据摄取、转换、编排、质量检查与元数据管理整合到单一平台详见 05-data-platforms/notes/01-introduction.md。本课用 DuckDB 作为本地托管数据库搭建三条分层流水线Ingestion layer摄取层从数据源抽取数据以原始格式落地存储Staging layer中间/清洗层预处理、清洗、转换并与查找表lookup table做关联Reports layer报表层对清洗后的数据进行聚合与计算产出报表。层与层之间通过 asset 的依赖关系dependencies串联Bruin 依据这些依赖构建出数据血缘data lineage图并据此决定编排执行顺序。这也是后续作业与血缘可视化的基础。项目初始化与目录结构从 Zoomcamp 模板初始化bruin init zoomcamp my-taxi-pipeline cd my-taxi-pipelinebruin init会从官方模板创建项目、初始化 git、生成.gitignore其中自动忽略.bruin.yml因为它包含数据库连接与密钥并创建项目配置文件。模板是一个基于 TODO 注释的练习工程你需要按内联注释逐步补全配置与代码本文即对应的完成版参考实现。完整初始化的其他细节可参见 05-data-platforms/notes/02-getting-started.md。初始化后的目录结构zoomcamp/ ├── .bruin.yml ├── README.md └── pipeline/ ├── pipeline.yml └── assets/ ├── ingestion/ │ ├── trips.py │ ├── requirements.txt │ ├── payment_lookup.asset.yml │ └── payment_lookup.csv ├── staging/ │ └── trips.sql └── reports/ └── trips_report.sql这个结构体现了 Bruin 的asset 命名约定asset 名称默认可以由文件路径推断习惯上按 schema/数据集分组——ingestion/下的文件映射为ingestion.*表staging/映射为staging.*reports/映射为reports.*。三个子目录天然对应三层架构后续每个 asset 通过depends声明依赖Bruin 据此绘制血缘并调度执行可参考 05-data-platforms/notes/06-core-01-projects.md 与 05-data-platforms/notes/06-core-03-assets.md。.bruin.yml环境与连接.bruin.yml位于项目根目录用于定义环境environment与连接connection且始终只保存在本地被.gitignore忽略绝不可提交到仓库因为它包含数据库路径、令牌等敏感信息。default_environment: default environments: default: connections: duckdb: - name: duckdb-default path: duckdb.dbdefault_environment指定默认使用的环境保证日常操作默认跑在开发环境避免误连生产environments.name.connections在每个环境下声明连接。duckdb连接只需给出name连接标识供 pipeline 与命令引用和path本地 DuckDB 数据库文件路径此处为duckdb.dbBruin 内置多种连接类型DuckDB、MotherDuck、PostgreSQL、MySQL、BigQuery、Redshift、Snowflake以及用于 API Key 等场景的自定义连接。更多环境与连接的高级用法如 production 环境配 BigQuery见 05-data-platforms/notes/06-core-01-projects.md。pipeline.yml管道定义pipeline.yml定义管道的名称、调度、默认连接与自定义变量name: nyc_taxi schedule: daily start_date: 2022-01-01 default_connections: duckdb: duckdb-default variables: taxi_types: type: array items: type: string default: [yellow]关键配置项配置项作用name管道标识schedule调度频率如hourly、daily、monthly或 cron 表达式每个管道只有一个调度start_date做全量刷新full refresh时从该日期开始处理数据default_connections指定本管道使用的连接连接定义在项目级.bruin.yml管道级做作用域限定避免密钥过度暴露见 05-data-platforms/notes/06-core-02-pipelines.mdvariables管道级自定义变量例如taxi_types用于控制摄取哪些出租车类型yellow、green 或两者兼顾支持array等类型并带 schema 校验变量可在运行时通过--var覆盖例如只跑绿色出租车bruin run ./pipeline/pipeline.yml --var taxi_types[green]start_date/end_date是 Bruin 的内置变量由调度周期决定如 daily 为当天起止、monthly 为整月起止可以在 SQL 中用 Jinja 模板{{ start_datetime }}注入、在 Python 中通过环境变量读取自定义变量在 Python 中则通过BRUIN_VARSJSON读取。变量机制的完整说明见 05-data-platforms/notes/06-core-04-variables.md。摄取层Ingestion摄取层包含三类 assetPython 脚本trips.py抓取出租车行程数据、Seed 文件payment_lookup.asset.ymlpayment_lookup.csv静态支付方式查找表、以及 Python 依赖清单requirements.txt。Python assettrips.pyPython asset 通过文件顶部bruin ... bruin装饰器声明元数据正文实现materialize()函数连接 NYC Taxi 开放数据 API 并抽取数据bruin name: ingestion.trips type: python image: python:3.11 materialization: type: table strategy: append columns: - name: pickup_datetime type: timestamp description: When the meter was engaged - name: dropoff_datetime type: timestamp description: When the meter was disengaged bruin import os import json import pandas as pd def materialize(): start_date os.environ[BRUIN_START_DATE] end_date os.environ[BRUIN_END_DATE] taxi_types json.loads(os.environ[BRUIN_VARS]).get(taxi_types, [yellow]) # Generate list of months between start and end dates # Fetch parquet files from: # https://d37ci6vzurychx.cloudfront.net/trip-data/{taxi_type}_tripdata_{year}-{month}.parquet return final_dataframe要点解读materialize()返回 DataFrameBruin 负责把返回的 DataFrame 写入目标数据库表你无需手写建表或 INSERTstrategy: append每次运行只追加新数据不动已有行适合原始数据持续落地的摄取场景BRUIN_START_DATE/BRUIN_END_DATEBruin 注入的环境变量标识本次运行的时间窗口代码据此生成需要拉取的月份列表BRUIN_VARS以 JSON 字符串注入的自定义变量集合这里通过json.loads(...).get(taxi_types, [yellow])读取管道变量taxi_types默认取[yellow]注释中给出的数据源 URL 即 NYC TLC 公开的 Parquet 文件托管地址格式为{taxi_type}_tripdata_{year}-{month}.parquet。Python asset 的通用写法与装饰器语法可对照 05-data-platforms/notes/06-core-03-assets.md变量在 Python 侧的完整访问方式含BRUIN_VAR_*前缀形式见 05-data-platforms/notes/06-core-04-variables.md。Seed 文件payment_lookup.asset.ymlSeed 资产用于把本地 CSV 文件直接导入数据库成为一张表适合支付方式这类静态参考数据name: ingestion.payment_lookup type: duckdb.seed parameters: path: payment_lookup.csv columns: - name: payment_type_id type: integer description: Numeric code for payment type primary_key: true checks: - name: not_null - name: unique - name: payment_type_name type: string description: Human-readable payment type checks: - name: not_null对应的payment_lookup.csv内容payment_type_id,payment_type_name 0,flex_fare 1,credit_card 2,cash 3,no_charge 4,dispute 5,unknown 6,voided_trip要点解读type: duckdb.seed表示目标连接类型为 DuckDB 的 seedparameters.path指向相对 asset 文件位置的 CSV列定义中除类型与描述外还支持primary_key与checks数据质量检查质量检查自动执行not_null非空、unique唯一等检查会在 asset 执行结束后自动运行失败即标记该资产失败从源头保障下游数据质量。requirements.txtPython 依赖pandas requests pyarrow python-dateutil这四者是摄取脚本运行所需pandas负责 DataFrame 处理requests用于 HTTP 请求pyarrow用于读写 Parquet 文件python-dateutil提供日期解析能力。Bruin 会接管运行环境在管道内本地安装这些依赖无需手工创建虚拟环境。清洗层StagingSQL assetstaging/trips.sql清洗层把原始行程数据与支付方式查找表关联做类型/字段投影与去重。SQL asset 同样以/* bruin ... bruin */注释块声明元数据/* bruin name: staging.trips type: duckdb.sql depends: - ingestion.trips - ingestion.payment_lookup materialization: type: table strategy: time_interval incremental_key: pickup_datetime time_granularity: timestamp columns: - name: pickup_datetime type: timestamp primary_key: true checks: - name: not_null custom_checks: - name: row_count_greater_than_zero query: | SELECT CASE WHEN COUNT(*) 0 THEN 1 ELSE 0 END FROM staging.trips value: 1 bruin */ SELECT t.pickup_datetime, t.dropoff_datetime, t.pickup_location_id, t.dropoff_location_id, t.fare_amount, t.taxi_type, p.payment_type_name FROM ingestion.trips t LEFT JOIN ingestion.payment_lookup p ON t.payment_type p.payment_type_id WHERE t.pickup_datetime {{ start_datetime }} AND t.pickup_datetime {{ end_datetime }} QUALIFY ROW_NUMBER() OVER ( PARTITION BY t.pickup_datetime, t.dropoff_datetime, t.pickup_location_id, t.dropoff_location_id, t.fare_amount ORDER BY t.pickup_datetime ) 1要点解读depends依赖声明同时依赖ingestion.trips与ingestion.payment_lookup保证本 asset 一定在两个摄取 asset 完成之后才运行strategy: time_intervalincremental_key: pickup_datetime这是 Bruin 基于时间列的增量物化策略——先删除当前时间窗口内的旧行再插入本次查询结果delete insert 的组合WHERE 必须过滤同一时间窗口查询里的{{ start_datetime }}/{{ end_datetime }}是 Jinja 模板变量由 Bruin 编译为实际日期且必须与物化的时间窗口保持一致否则会造成重复数据QUALIFY ROW_NUMBER()去重按pickup_datetime, dropoff_datetime, pickup_location_id, dropoff_location_id, fare_amount组成的复合键编号只保留每组第 1 行剔除重复行程custom_checks自定义检查除了列级checks还支持针对整表的 SQL 断言——这里用一条查询判断staging.trips行数是否大于 0结果为 1 才通过。报表层ReportsSQL assetreports/trips_report.sql报表层对清洗后的数据做聚合产出按天、按车型、按支付方式统计的报表表/* bruin name: reports.trips_report type: duckdb.sql depends: - staging.trips materialization: type: table strategy: time_interval incremental_key: trip_date time_granularity: date columns: - name: trip_date type: date primary_key: true - name: taxi_type type: string primary_key: true - name: payment_type type: string primary_key: true - name: trip_count type: bigint checks: - name: non_negative bruin */ SELECT CAST(pickup_datetime AS DATE) AS trip_date, taxi_type, payment_type_name AS payment_type, COUNT(*) AS trip_count, SUM(fare_amount) AS total_fare, AVG(fare_amount) AS avg_fare FROM staging.trips WHERE pickup_datetime {{ start_datetime }} AND pickup_datetime {{ end_datetime }} GROUP BY 1, 2, 3要点解读同样使用time_interval增量策略但incremental_key变为trip_date、time_granularity变为date与报表粒度的日期列对齐输出列带列级质量检查trip_count声明non_negative检查确保聚合结果非负聚合逻辑COUNT(*)统计行程数、SUM(fare_amount)汇总车费、AVG(fare_amount)计算平均车费按(trip_date, taxi_type, payment_type)三维分组WHERE同样按时间窗口过滤与增量物化窗口保持一致。运行整条管道核心 CLI 命令# Validate structure and definitions bruin validate ./pipeline/pipeline.yml # Run with a small date range for testing bruin run ./pipeline/pipeline.yml --start-date 2022-01-01 --end-date 2022-02-01 # Full refresh bruin run ./pipeline/pipeline.yml --full-refresh # Query results bruin query --connection duckdb-default --query SELECT COUNT(*) FROM ingestion.trips命令语义命令作用bruin validate path只校验不执行检查 asset 定义是否正确、连接是否配置、血缘中是否存在循环依赖、有无断裂引用运行前务必先 validatebruin run path --start-date ... --end-date ...在指定日期范围内执行管道测试阶段建议先用小日期范围如上例仅 1 个月bruin run --full-refresh全量刷新删表并从头重建常用于首次运行或 schema 变更后bruin run --var ...运行时覆盖管道变量如--var taxi_types[yellow,green]bruin run --asset X --downstream/--upstream只运行某个 asset 及其下游/上游依赖bruin query --connection conn --query ...对指定连接执行即席 SQL 查询bruin lineage path查看 asset 依赖血缘图完整命令参考见 05-data-platforms/notes/06-core-05-commands.md。执行顺序与血缘可视化打开 Bruin 面板VS Code / Cursor 扩展中的管道 YAML 文件并切换到lineage血缘标签页即可看到所有 asset 及其依赖关系。本管道的执行顺序为摄取 asset 先运行ingestion.trips与ingestion.payment_lookup二者并行清洗 assetstaging.trips在两个摄取 asset 全部完成后运行报表 assetreports.trips_report在清洗完成后运行。这正是依赖声明带来的编排效果Bruin 从depends构建血缘图并做拓扑排序保证任何依赖先于被依赖方执行。若使用 Bruin MCP 与 AI Agent 协作可以让 Agent 自动完成创建资产、配置物化策略、设置质量检查、bruin validate、测试日期运行与bruin query验证的完整闭环具体见 05-data-platforms/notes/04-bruin-mcp.md。物化策略总结Bruin 为 SQL / Python asset 提供多种物化策略本课涉及的核心策略如下策略行为table每次运行删除并重建整张表append只插入新数据不触碰已有行merge基于键列做 Upsert存在则更新、不存在则插入time_interval删除日期范围内的行再重新插入该范围数据deleteinsert删除匹配行后插入createreplace创建或替换整张表选择建议原始数据落地用append按时间窗口增量刷新用time_interval如本课的 staging 与 reports 层键维度更新用merge一次性重建用table或createreplace。也可参见 05-data-platforms/notes/06-core-03-assets.md 中的策略对比。延伸阅读与配套练习入门与环境安装05-data-platforms/notes/02-getting-started.mdBruin 平台定位与数据栈概览05-data-platforms/notes/01-introduction.md项目 / 管道 / 资产 / 变量 / 命令核心概念06-core-01-projects.md、06-core-02-pipelines.md、06-core-03-assets.md、06-core-04-variables.md、06-core-05-commands.md用 AI Agent 构建同款管道05-data-platforms/notes/04-bruin-mcp.md配套作业含管道结构、物化策略、变量覆盖、依赖运行、质量检查、血缘与全量刷新等 7 道题cohorts/2026/05-data-platforms/homework.md通过本课你已经用 Bruin 在本地 DuckDB 上跑通了一条覆盖原始数据摄取 → 清洗关联去重 → 聚合报表三层的完整 ELT 管道配置环境与连接、声明依赖与血缘、选择增量物化策略、内建质量检查并用命令行完成校验、运行与查询。后续可以在此基础上把同一套资产定义迁移到 BigQuery、Snowflake 等云端连接或结合 Bruin MCP 用自然语言驱动 Agent 迭代管道向生产级数据平台演进。【免费下载链接】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),仅供参考
返回列表