dlt 数据管道核心术语详解:Source、Resource、Pipeline 与 Schema 的完整概念图谱
【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy 🛠️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt
dlt 的官方文档以一份 Glossary(术语表)定义了整套数据加载体系的八个核心概念:Source、Resource、Destination、Pipeline、Verified Source、Schema、Config 与 Credentials。本文以这份术语表为骨架逐一展开,并结合 dlt 仓库中的源码实现(source.py、resource.py)说明每个术语背后的真实机制,帮助你在编写管道代码时准确使用每一对概念。
Source:数据的逻辑分组与提取入口
按术语表的定义,Source是"持有一定结构数据的场所,组织为一个或多个资源(resource)"。文档用三个类比把它讲得很直白:
- 如果 API 的端点是资源,那么 API 就是数据源;
- 如果电子表格的标签页是资源,那么电子表格就是数据源;
- 如果数据库的表是资源,那么数据库就是数据源。
术语表特别强调:在 dlt 文档体系中,source 同时指代软件组件——即一个用于从数据场所提取数据、由一个或多个 resource 组件构成的 Python 函数。
源码视角:DltSource 类
在 dlt 源码中,被@dlt.source装饰的函数调用后会返回 DltSource 实例。其类 docstring 说明了它自动化了以下能力:
- 可直接传给
pipeline.run()以加载其中全部 resource 的数据; - 通过
with_resources()方法选择/取消选择要加载的 resource; - 每个 resource(
DltResource实例)都可以作为 source 的属性直接访问; - 实现了
Iterable接口,即使没有 pipeline 也能自行遍历数据; - 会为 resource 和 transformer 构建 DAG(有向无环图),并优化提取,使父 resource 只被提取一次;
- 可获取 source 及其所有 resource 的 schema;
- 提供
run()方法,用 dlt 的默认 pipeline 实例直接加载数据。
从源码结构看,DltSource内部用一个 DltResourceDict 字典统一管理 resource:其selected属性返回"将被提取并加载到目的地的 resource 子集",extracted属性则返回"所有将被提取的 resource(包含所选 resource 及其父级)"。这解释了为什么pipeline.run(source.with_resources("companies", "deals"))可以只加载部分端点——选择操作本质上是在修改字典中每个 resource 的selected标志。
DltSource还实现了几个在文档中反复出现的便捷属性,源码中可直接验证:
- max_table_nesting:设置嵌套表的最大深度,超出深度的部分以 struct 或 JSON 加载。源码显示它最终写入 schema 的 normalizer 配置(
x-normalizer的max_nesting键)。max_table_nesting=0不生成任何嵌套表,1只生成根表的嵌套表,文档建议对 MongoDB 这类深度嵌套数据源设为 2 或 3 以获得最清晰的 schema; - root_key:把根表的
_dlt_id外键传播到所有嵌套表,在将 resource 的写模式切换为 merge 时特别有用; - schema_contract:读取或设置 source 级 schema 合同。
编写 source 时的关键约定(来自 source 概念文档):source 函数在调用时立即执行,而 resource 像 Python 生成器一样延迟执行。因此不应在 source 函数体内做提取操作(如反射数据库表),而应把这些工作留给 resource——这样可以在pipeline.run或pipeline.extract阶段获得错误处理、执行指标和并行化等收益。
Resource:数据的最小提取单元
术语表将Resource定义为"数据源内数据的一个逻辑分组,通常持有结构和来源相似的数据",并给出与 Source 对应的一组类比:API 源中的 resource 是端点,电子表格源中的 resource 是标签页,数据库源中的 resource 是表。同样地,resource 也指代软件组件——一个从数据源场所提取数据的 Python 函数。
源码视角:DltResource 与数据管道
在 resource.py 中,@dlt.resource装饰器把被装饰函数包装为DltResource。从源码可以看到,每个 resource 内部持有一条"管道"(pipe),各种变换操作都是向这条管道追加步骤:
- add_map:追加
MapItem步骤,逐条转换数据项; - add_filter:追加
FilterItem步骤,返回True的数据项被保留; - add_yield_map:追加
YieldMapItem步骤,一个数据项可派生 0 个或多个新数据项(行转列/展开); - add_metrics:追加
MetricsItem步骤,在修改数据项的同时收集自定义指标; - add_limit:按"yield 次数"(
max_items)或时间(max_time)限制提取量。
所有这些方法都返回self,因此可以链式调用,例如:
users().add_metrics(track_filtered).add_filter(lambda u: u["user_id"] != "me").add_map(anonymize_user)从源码实现可以确认文档中的一个重要细节:add_limit限制的是yield 的次数而非行数。每次 yield 可能包含一个列表(一页数据),达到限制后 dlt 会关闭产生数据的迭代器/生成器。若希望按行数计数,可传count_rows=True。
resource 常用的装饰器参数(见 resource 概念文档)包括:
name:生成的表名,默认为被装饰函数名;write_disposition:加载方式,支持append、replace、merge,默认append;table_name、primary_key/merge_key、columns(列的TTableSchemaColumns类型提示,如把tags列声明为json类型避免拆分为嵌套表);nested_hints:定义嵌套表 schema,深层嵌套可用元组路径定位,如("purchases", "coupons");parallelized=True:同步 resource 的并行提取;async 生成器则自动并发提取。
此外,resource 还内置了 max_table_nesting 属性(同样落在x-normalizer提示中,未设置时回退到 source 级或默认值 1000),以及 select_tables 方法——对动态分发到多张表的 resource(如按event["type"]分表的资源流),可用它筛选接收数据的表。
Destination:数据最终落地的存储
术语表中Destination的定义简洁明确:"源数据被加载到的数据存储(例如 Google BigQuery)"。
在仓库中,目的地实现位于 dlt/destinations/impl/ 目录,包含 bigquery、clickhouse、duckdb、postgres、snowflake、databricks、filesystem、lance、qdrant、weaviate 等 20 余种内置实现的子包;目的地能力、客户端与作业抽象则由 dlt/destinations/ 下的job_client_impl.py、sql_jobs.py、insert_job_client.py等文件提供。加载到某个目的地时,pipeline 会把 normalize 后的数据写成加载包(load package),再由对应目的地的作业客户端执行插入、合并或替换。
Pipeline:连接 Source 与 Destination 的执行者
术语表对Pipeline的定义是:"按照 schema 提供的指令,把数据从 source 移动到 destination(即执行提取、规范化、加载)"。
这是 dlt 三大执行阶段的载体:extract(从 source 的 resource 迭代数据)、normalize(推断/应用 schema、生成加载包)、load(把文件写入 destination)。实现入口在 dlt/pipeline/pipeline.py,核心方法pipeline.run(source_or_resources)接受单个 source、多个 source、单个 resource 或它们的列表;默认所有传入的 source 会加载到同一个 dataset,也可以把一个 source 拆解为多个 pipeline(例如把 50 张表的复制作业拆成并行度更高的多个 DAG 任务)。
术语表中四个概念在 pipeline 运行时串联为:Config/Credentials在运行时注入 →Source的 resource 被Pipeline提取 →Schema描述规范化后的数据并指导加载 → 数据写入Destination。
Verified Source:随 dlt 分发的经验证数据源模块
Verified source是一个随dlt init分发的 Python 模块,允许你创建从特定 Source 提取数据的 pipeline。这类模块旨在公开发布,供他人用它来构建 pipeline。
术语表给出了"verified(已验证)"的完整判据:一个 source 必须经过发布才能成为 verified,这意味着它具备:
- 测试(tests);
- 测试数据(test data);
- 演示脚本(demonstration scripts);
- 文档(documentation);
- 其生成的数据集经过数据工程师审阅(reviewed by a data engineer)。
仓库中这类内置源位于 dlt/sources/(如sql_database、filesystem、rest_api、config、credentials),由 dlt/sources/init.py 导出DltSource、DltResource、incremental等构建块。其"可被其他 source 复用与改名"的能力在源码中有对应实现:SourceReference/AnySourceFactory支持source.clone(name=..., section=...)生成改名副本,从而把配置放入自定义配置段或让多个同名源实例并存。
Schema:规范化数据的结构描述与加载指令
术语表对Schema的定义包含两层含义:
- 描述结构:规范化后数据的结构(例如展开后的表、列类型等);
- 提供指令:数据应如何被处理和加载——即告诉 dlt 数据的内容是什么、以及如何将其加载到 destination。
在仓库中,schema 的核心类是 dlt/common/schema/schema.py 中的Schema,类型定义(表 schema、列 schema、合同设置等)位于 dlt/common/schema/typing.py,类型探测逻辑在 dlt/common/schema/detections.py。DltSource.schema属性让你在提取前即可检查并修改 schema(加表、改列定义等),而 resource 的compute_table_schema()方法可以输出该 resource 将生成的表 schema(动态提示可传入一个示例数据项来求值)。
值得注意的实现细节:前面提到 source 级的max_table_nesting和root_key并不是独立配置项,而是写入 schema 的 normalizer 配置中(见 source.py 的 setter 实现)。从源码结构看,这印证了术语表的第二层定义——schema 不只是"描述",它携带着影响规范化行为(嵌套展开深度、外键传播)的执行指令。
Config:运行时注入的行为配置
术语表将Config定义为"运行时传递给 pipeline 的一组值(例如用于在本地与生产环境中改变其行为)"。
dlt 的配置系统位于 dlt/common/configuration/ 目录:container.py提供依赖注入容器,inject.py实现向函数参数的自动注入,specs/子目录存放各类配置规格定义。实践中的典型用法是在资源函数参数上标注:
@dlt.resource def fs_resource(bucket_url=dlt.config.value): ... pipeline.run(fs_resource("s3://my-bucket/reports"), table_name="reports")运行时,dlt 按配置段的查找顺序(命令行、环境变量、config.toml等)解析值并自动注入。配置文件的写入与查找规则在 dlt/_workspace/config_toml_writer.py 与配置文档(setup 指南 中 "How dlt looks for values" 一节)中有详细说明。改名 source 的场景中,Config 还可使用紧凑布局,如[sources.my_db],完整路径sources.my_db.my_db同时存在时优先。
Credentials:永不以明文共享的配置子集
术语表对Credentials的定义是:"配置的一个子集,其元素被保密保存,绝不以明文形式共享"。
实现上,Credentials 在 dlt/common/configuration/specs/ 中有专门的凭据规格(如CredentialsSpecification及其子类),敏感值通常配合secrets.toml存放,环境变量中的密钥则用dlt.secrets.value注解声明:
@dlt.source def hubspot(api_key=dlt.secrets.value): ...Config 与 Credentials 的分工可以概括为:Config 管"行为"(路径、表名、批量大小等非敏感参数),Credentials 管"身份"(API key、密码、连接串等敏感参数),二者都通过同一套依赖注入机制在运行时解析到函数参数中。
八个术语在一条 pipeline 中的位置
把术语表串联起来,一条 dlt 管道的完整生命周期是:
import dlt # 1. Source:声明提取入口,参数用 Credentials/Config 注解 @dlt.source def hubspot(api_key=dlt.secrets.value): # 2. Resource:每个端点一个 resource(延迟执行) @dlt.resource(name="companies") def companies(): yield requests.get(BASE + "/companies", headers=h).json() yield companies # 3. Pipeline:按 4. Schema 的指令执行 extract → normalize → load pipeline = dlt.pipeline(pipeline_name="hubspot", destination="duckdb") info = pipeline.run(hubspot())其中:hubspot()返回DltSource,其内部以DltResourceDict管理 resource;运行时 dlt 从 Config/Credentials 系统解析api_key;提取数据时由 schema 推断列类型并生成嵌套表;最终由 destination 对应的作业客户端完成加载。理解了这套术语与 dlt/extract/ 源码中对应类的映射关系,你就可以在调试管道(检查source.resources.selected、查看compute_table_schema()输出)时准确定位问题所在的概念层级。
【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy 🛠️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考