dlt 数据管道核心术语详解:Source、Resource、Pipeline 与 Schema 的完整概念图谱
发布时间:2026/9/17 2:24:32 作者:尧图编辑部 阅读量:1,286

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/dltdlt 的官方文档以一份 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每个 resourceDltResource实例都可以作为 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_nesting0不生成任何嵌套表1只生成根表的嵌套表文档建议对 MongoDB 这类深度嵌套数据源设为 2 或 3 以获得最清晰的 schemaroot_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_rowsTrue。resource 常用的装饰器参数见 resource 概念文档包括name生成的表名默认为被装饰函数名write_disposition加载方式支持append、replace、merge默认appendtable_name、primary_key/merge_key、columns列的TTableSchemaColumns类型提示如把tags列声明为json类型避免拆分为嵌套表nested_hints定义嵌套表 schema深层嵌套可用元组路径定位如(purchases, coupons)parallelizedTrue同步 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_urldlt.config.value): ... pipeline.run(fs_resource(s3://my-bucket/reports), table_namereports)运行时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_keydlt.secrets.value): ...Config 与 Credentials 的分工可以概括为Config 管行为路径、表名、批量大小等非敏感参数Credentials 管身份API key、密码、连接串等敏感参数二者都通过同一套依赖注入机制在运行时解析到函数参数中。八个术语在一条 pipeline 中的位置把术语表串联起来一条 dlt 管道的完整生命周期是import dlt # 1. Source声明提取入口参数用 Credentials/Config 注解 dlt.source def hubspot(api_keydlt.secrets.value): # 2. Resource每个端点一个 resource延迟执行 dlt.resource(namecompanies) def companies(): yield requests.get(BASE /companies, headersh).json() yield companies # 3. Pipeline按 4. Schema 的指令执行 extract → normalize → load pipeline dlt.pipeline(pipeline_namehubspot, destinationduckdb) 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),仅供参考