Data Engineering Zoomcamp 2026 实战:用 dlt + MCP 从零构建纽约出租车数据管道
【免费下载链接】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 2026 dlt 工作坊的课后作业(对应仓库文档 cohorts/2026/workshops/dlt/dlt_homework.md)为主线,完整讲解如何从一个没有现成 dlt 脚手架的自定义 REST API 出发,借助 dlt(data load tool)与 dlt MCP Server 构建生产级数据管道,将分页 JSON 数据加载到本地 DuckDB,并通过 dlt Dashboard、dlt MCP 对话与 marimo Notebook 三种方式完成数据探索与作业问题解答。读完本文,你将掌握 dlt 项目初始化、rest_api_source配置、分页处理、AI 辅助开发工作流,以及多种数据检视方法,可将其直接复用到你自己的 API 数据管道项目中。
一、作业背景与目标
本次作业源于 Data Engineering Zoomcamp 2026 的 dlt 工作坊。工作坊主体(见 workshops/dlt/README.md)演示了如何基于 dlt 官方脚手架(scaffold)快速搭建 Open Library API 管道,并利用 AI 辅助 IDE 生成、调试和运行管道。而本次作业将场景升级为从零构建:数据源是一个没有现成 dlt 脚手架的自定义 API,你必须自行提供 API 元数据并驱动 AI 助手生成管道代码。
作业的核心目标:
- 构建一个可运行的 dlt 管道,从自定义 API 提取纽约出租车行程数据;
- 将数据加载到本地 DuckDB 数据库(无需任何云凭据或配置);
- 通过 dlt Dashboard、dlt MCP Server 与 marimo Notebook 三种方式探索加载后的数据;
- 回答基于数据的 3 个分析问题。
本次作业的独特价值在于:它不仅训练传统的管道构建技能,还引入了AI 辅助开发(AI-assisted development)工作流——dlt MCP Server 让 AI 助手能够访问 dlt 官方文档、代码示例与管道元数据,从而实现"用自然语言描述 API → 生成管道代码 → 自动运行加载"的端到端体验。
二、数据源:NYC Yellow Taxi 自定义 API
作业使用的数据源是NYC 黄色出租车行程数据,通过一个自定义 REST API 提供。该 API 的特点如下表:
| 属性 | 值 |
|---|---|
| 基础 URL | https://us-central1-dlthub-analytics.cloudfunctions.net/data_engineering_zoomcamp_api |
| 数据格式 | 分页 JSON(Paginated JSON) |
| 每页大小 | 1,000 条记录/页 |
| 分页终止条件 | 当返回空页时停止 |
这个 API 是本次作业的关键难点:
- 没有现成脚手架:
dlt init只提供 dlt 项目骨架,不会生成 API 元数据 YAML 文件,因此 API 的 URL、分页方式等全部需要你(或你的 AI 助手)自行配置。 - 分页必须正确处理:API 按每页 1,000 条返回,分页终止条件是"返回空页",即当某页没有任何记录时,表示数据已全部拉取完毕。如果分页配置错误,会导致数据不完整或无限循环。
- 数据内容:纽约出租车行程记录,包含行程时间、支付方式、金额、小费金额等字段,用于回答作业中的三个分析问题。
三、准备工作:环境与工具链
在开始构建管道之前,需要准备以下环境(与工作坊 README 保持一致,见 workshops/dlt/README.md):
1. Python 3.11+
python --version # 应显示 3.11 或更高版本2. uv 或 pip
工作坊推荐使用 uv 管理依赖:
curl -LsSf https://astral.sh/uv/install.sh | sh3. Agentic IDE
需要一个具备 AI 辅助能力的代码编辑器,工作坊推荐:
| IDE | 说明 |
|---|---|
| Cursor | VS Code 分支,内置 AI 辅助(工作坊推荐) |
| Windsurf | 备选 agentic IDE |
| VS Code + GitHub Copilot | 可用,但集成度略低 |
4. 理解 dlt 基础概念(可选但推荐)
如果你对 dlt 还不熟悉,工作坊提供了概念讲解 Notebook(对应仓库文件 dlt_Pipeline_Overview.ipynb),其中核心概念如下:
- Source(数据源):管道中负责从某处获取数据的部分,在 dlt 中通常用
rest_api_source以简单字典配置描述 API,而非手写大量请求代码; - Pipeline(管道):描述数据的目标位置(如 DuckDB),并跟踪表、schema 与运行历史;
- Extract → Normalize → Load 三阶段:这是 dlt 管道的完整流程,详见下文第五节。
四、环境配置:搭建 dlt 开发环境
Step 1:创建项目文件夹
如果你在工作坊演示中已经创建过项目文件夹(如 Open Library 演示项目),可以直接复用;否则新建:
mkdir taxi-pipeline cd taxi-pipeline然后在 Cursor(或你偏好的 agentic IDE)中打开该文件夹。
Step 2:配置 dlt MCP Server(如未配置)
dlt MCP Server 是本次 AI 辅助开发的核心组件。它让 AI 助手能够访问 dlt 文档、代码示例以及你的管道元数据(pipeline metadata),从而更准确地生成和调试代码。
Cursor 配置方式:进入Settings → Tools & MCP → New MCP Server,添加:
{ "mcpServers": { "dlt": { "command": "uv", "args": [ "run", "--with", "dlt[duckdb]", "--with", "dlt-mcp[search]", "python", "-m", "dlt_mcp" ] } } }VS Code(Copilot)配置方式:在项目文件夹中创建.vscode/mcp.json:
{ "servers": { "dlt": { "command": "uv", "args": [ "run", "--with", "dlt[duckdb]", "--with", "dlt-mcp[search]", "python", "-m", "dlt_mcp" ] } } }Claude Code 配置方式:在终端中运行:
claude mcp add dlt -- uv run --with "dlt[duckdb]" --with "dlt-mcp[search]" python -m dlt_mcp配置要点说明:三条命令本质相同——使用
uv run临时创建运行环境,安装dlt[duckdb](dlt 主库 + DuckDB 目标支持)与dlt-mcp[search](MCP Server 及其文档搜索能力),然后以 Python 模块方式启动dlt_mcp。启动后,AI 助手即可通过 MCP 协议查询 dlt 文档、代码示例与当前管道元数据。
Step 3:安装 dlt
pip install "dlt[workspace]"dlt[workspace]是一个聚合安装,包含 dlt 核心、常用目标数据库支持,以及 MCP Server 等 AI 辅助开发工具。
Step 4:初始化项目
dlt init dlthub:taxi_pipeline duckdb该命令会:
- 创建 dlt 项目文件(如
.dlt/配置目录、requirements.txt等); - 创建用于 AI 辅助的 Cursor 规则(Cursor rules);
- 但不会创建 API 元数据 YAML 文件——因为
dlthub:taxi_pipeline没有对应的脚手架。
这也是本次作业与工作坊演示(Open Library 有脚手架)的核心区别:你需要在下一步自行提供 API 信息。
五、从零构建管道:Extract → Normalize → Load
由于该 API 没有脚手架,你需要将 API 详细信息写进提示词(prompt),让 AI 助手生成管道代码。
Step 5:用提示词驱动 Agent 生成管道
以下是作业给出的示例提示词(可复制到 Cursor 等 agentic IDE 的对话中):
Build a REST API source for NYC taxi data. API details: - Base URL: https://us-central1-dlthub-analytics.cloudfunctions.net/data_engineering_zoomcamp_api - Data format: Paginated JSON (1,000 records per page) - Pagination: Stop when an empty page is returned Place the code in taxi_pipeline.py and name the pipeline taxi_pipeline. Use @dlt rest api as a tutorial.提示词的关键要素:
- 明确说明是自定义 REST API,无脚手架可用;
- 提供完整 API 元数据:Base URL、数据格式(分页 JSON)、分页规则(每页 1,000 条、空页终止);
- 指定输出文件与管道命名:
taxi_pipeline.py、管道名taxi_pipeline; - 引用
@dlt rest api教程,让 AI 采用 dlt 官方的 REST API 最佳实践。
理解 Agent 将生成的代码
为了让你知道 agent 会生成什么、并能验证其正确性,这里给出工作坊演示中 Open Library 管道的完整代码(见仓库 open_library_pipeline.py),它展示了 dlt REST API 源的标准写法:
"""Pipeline to ingest data from the Open Library Search API.""" import dlt from dlt.sources.rest_api import rest_api_source def open_library_source(query: str = "harry potter"): """ Create a dlt source for the Open Library Search API. Args: query: Search query string (default: "harry potter") """ return rest_api_source({ "client": { "base_url": "https://openlibrary.org", }, "resource_defaults": { "primary_key": "key", "write_disposition": "replace", }, "resources": [ { "name": "books", "endpoint": { "path": "search.json", "params": { "q": query, "limit": 100, }, "data_selector": "docs", "paginator": { "type": "offset", "limit": 100, "offset_param": "offset", "limit_param": "limit", "total_path": "numFound", }, }, }, ], }) if __name__ == "__main__": pipeline = dlt.pipeline( pipeline_name="open_library_pipeline", destination="duckdb", dataset_name="open_library_data", progress="log", ) # Load Harry Potter books from Open Library load_info = pipeline.run(open_library_source(query="harry potter")) print(load_info)对照该示例,你的出租车管道在生成时应重点关注以下几个配置点:
| 配置项 | 作用 | 出租车作业中的注意点 |
|---|---|---|
client.base_url | API 基础 URL | 应替换为作业给定的自定义 API URL |
resources[].name | 资源名称 | 决定生成表名 |
endpoint.path | API 路径 | 自定义 API 通常为根路径或固定路径 |
data_selector | 从 JSON 响应中选取数据数组 | 需根据 API 返回结构确定(如docs、results或根数组) |
paginator | 分页器配置 | 作业 API 为"空页终止"分页,可使用 dlt 的分页器处理;Open Library 示例用的是offset分页,逻辑可参考但不完全相同 |
write_disposition | 写入策略 | replace表示每次全量替换,适合此类作业场景 |
关于分页的关键提示:作业 API 的分页终止条件是"返回空页"(stop when an empty page is returned)。dlt 的 REST API 源内置多种分页器(offset、cursor、page number 等),agent 会为自定义 API 选择合适的配置。如果管道运行后加载的行数明显偏少(如只有 1,000 条),说明分页可能只拉取了第一页,需要检查 paginator 配置。
Step 6:运行与调试
生成代码后,运行管道:
python taxi_pipeline.py如果出现错误,将错误信息粘贴到对话中让 agent 调试。这是 AI 辅助开发的典型迭代流程:运行 → 报错 → 反馈给 AI → 修复 → 重跑。
管道运行成功后,dlt 会依次执行三个阶段:
- Extract(提取):向自定义 API 发送请求,下载原始 JSON 响应,存放到 dlt 本地工作目录。此时数据尚未进入 DuckDB。
- Normalize(归一化):将嵌套 JSON 转换为关系型表结构。dlt 会为每个表添加
_dlt_id(行唯一标识)与_dlt_load_id(关联加载任务)跟踪列,将嵌套列表展开为子表(如trips__xxx形式的子表,通过_dlt_parent_id关联父表),并创建_dlt_loads、_dlt_pipeline_state、_dlt_version等元数据表。 - Load(加载):在 DuckDB 中创建表(若不存在)并插入归一化后的数据,同时记录加载历史。
在完整理解这三阶段后,也可以直接用pipeline.run(source)一条命令完成全部流程(它等价于extract → normalize → load三步)。
六、探索加载后的数据:三种方法
管道运行成功后,作业要求使用工作坊中讲解的方法来调查数据。仓库中的分析示例(见 analysis.py)展示了 marimo 与 ibis 的组合用法,可作参考。
方法一:dlt Dashboard
dlt pipeline taxi_pipeline show这会启动一个 Web 应用,用于:
- 查看管道状态与运行历史;
- 浏览 schema、表与列结构;
- 查询已加载数据;
- 调试可能存在的问题。
方法二:dlt MCP Server 对话式查询
配置好 dlt MCP Server 后,可以直接在对话中向 AI 提问,例如:
"What tables were created in the pipeline?" "Show me the schema for the trips table." "How many rows were loaded?"
agent 可以访问你的管道元数据,直接回答这些问题,无需手写 SQL。
方法三:marimo Notebook 可视化分析
创建 marimo notebook 进行查询与可视化。工作坊的推荐运行方式:
marimo edit your_notebook.py # 编辑模式(开发) marimo run your_notebook.py # 运行模式(查看报告)工作坊的示例 notebook(analysis.py)展示了 dlt + marimo + ibis 的标准分析流程:
import marimo as mo import dlt import ibis import altair as alt from dlt.helpers.marimo import render, load_package_viewer # 使用 dlt 原生接口访问管道与数据集 pipeline = dlt.attach("open_library_pipeline") dataset = pipeline.dataset() # 获取 ibis 连接进行丰富的数据探索 ibis_con = dataset.ibis()核心步骤:
- 使用
dlt.attach("pipeline_name")重新挂载已存在的管道; - 通过
pipeline.dataset()获取数据集接口,无需手写 SQL; - 使用
dataset.ibis()获得 ibis 连接,进行链式查询(group_by、agg、order_by 等); - 使用 Altair 绘制图表;
- 通过
render(load_package_viewer)直接在 notebook 中嵌入 dlt 包查看器。
依赖清单
工作坊项目的依赖配置(见 pyproject.toml)可作为环境参考:
dependencies = [ "altair>=6.0.0", "dlt[workspace]>=1.21.0", "ibis-framework[duckdb]>=12.0.0", "jupyterlab>=4.5.4", "marimo>=0.19.9", ]七、作业问题与解题思路
管道成功运行后,请基于加载到 DuckDB 的数据回答以下三个问题:
Question 1:数据集的开始日期和结束日期是什么?
- 2009-01-01 至 2009-01-31
- 2009-06-01 至 2009-07-01
- 2024-01-01 至 2024-02-01
- 2024-06-01 至 2024-07-01
解题思路:对行程时间字段执行 MIN/MAX 聚合。注意行程时间字段可能是字符串格式,必要时使用 CAST 转换为日期类型后再取极值。
Question 2:使用信用卡支付的行程占比是多少?
- 16.66%
- 26.66%
- 36.66%
- 46.66%
解题思路:支付方式字段(如payment_type)中识别信用卡对应的枚举值,计算信用卡行程数 / 总行程数。建议用 GROUP BY 先查看该字段的取值分布,确认信用卡的取值后再计算比例。
Question 3:小费(tips)产生的总金额是多少?
- $4,063.41
- $6,063.41
- $8,063.41
- $10,063.41
解题思路:对小费金额字段执行 SUM 聚合。同样注意字段的数值类型转换。
提交与截止日期:通过课程官网的作业提交表单提交答案(对应
courses.datatalks.club/de-zoomcamp-2026/homework/dlt),注意截止时间以网站公布为准。
八、实用技巧与常见问题
作业文档给出了几条关键提示,值得展开说明:
API 返回分页数据,确保管道正确处理分页:这是最容易出错的地方。验证方法:对比加载的总行数与 API 返回的总记录数。若只有 1 页数据被加载(1,000 行),说明分页器配置不正确。由于终止条件是"空页",要确认 agent 生成的分页器在遇到空页时能正常结束循环,而不是报错或死循环。
Agent 卡住时,把错误信息粘贴到对话中让它调试:AI 辅助开发的核心理念就是"出错即反馈"。错误信息(traceback)是 agent 调试的最重要线索,粘贴完整报错,让它修复后重跑。
使用 dlt MCP Server 查询管道元数据:如"加载了多少行"、"建了哪些表"这类问题,直接问 agent 即可,无需手写 SQL。这既能验证管道正确性,也能练习 MCP 工作流。
尝试多种调查方法:作业鼓励在回答问题时尝试 dlt Dashboard、MCP 对话、marimo 可视化等不同方法,并分享哪种方法效果最好——这也是本次作业的隐性训练目标。
九、参考资料
| 资源 | 说明 |
|---|---|
| workshops/dlt/README.md | 工作坊完整指南(Open Library 演示管道) |
| dlt_Pipeline_Overview.ipynb | dlt 概念讲解 Notebook(Extract/Normalize/Load) |
| open_library_pipeline.py | 工作坊演示管道的完整代码(REST API 源标准写法) |
| analysis.py | marimo + ibis 数据探索示例 |
| pyproject.toml | 工作坊依赖清单 |
| images/etl_diagram.png | Extract → Normalize → Load 流程示意图 |
十、学习在公开中分享(Learning in Public)
课程鼓励所有学员公开分享学习成果("learning in public")。作业文档提供了 LinkedIn 与 Twitter/X 的发帖模板,其要点如下:
- 总结你在工作坊中学到的技能:REST API 数据管道、dlt MCP Server 的 AI 辅助开发、分页 API 数据加载到 DuckDB、dlt Dashboard 与 marimo 数据检视;
- 分享你的作业解决方案链接;
- 介绍 Data Engineering Zoomcamp 这一免费课程。
分享不仅是社区文化的体现,也是数据工程师建立个人品牌、沉淀知识体系的有效方式。
【免费下载链接】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),仅供参考