
dlt 目的端表结构详解数据集组织、嵌套表引用与 _dlt 内部管理表【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt本文基于 dlt 官方文档 Destination tables lineage系统讲解 dlt pipeline 运行后在目的端数据库中创建的各类表的结构与组织方式数据库 schema 与数据集命名、resource 到表的映射、嵌套数据如何分裂为带引用关系的子表、列/表命名归一化规则、variant 列机制、load package 与load_id追踪、merge 写入下的 staging 数据集、dev_mode版本化数据集以及_dlt_loads、_dlt_pipeline_state、_dlt_version三张内部管理表的列定义与用途。读完后你可以准确解释任意 dlt 目的端数据库里的表结构并能利用内置元数据完成数据血缘追溯、未完成加载过滤和增量加载排查。从一个 pipeline 开始目的端会创建什么先运行一个最简单的 dlt pipeline官方文档示例import dlt data [ {id: 1, name: Alice}, {id: 2, name: Bob} ] pipeline dlt.pipeline( pipeline_namequick_start, destinationduckdb, dataset_namemydata ) load_info pipeline.run(data, table_nameusers)运行后dlt 会在目的端数据库这里是 DuckDB 内存数据库中创建一个数据库 schema以及其中一张名为users的表并把你 source 中的数据写入其中。其他数据库目的端的行为与概念基本一致。小技巧可以使用 dlt pipeline CLI 的show命令查看目的端数据库中的表例如dlt pipeline showdashboard 文档中也有相关说明。数据库 schema 的命名数据库 schema 是承载你已加载数据的一组表。schema 名与 pipeline 定义中提供的dataset_name相同本例中显式设置了dataset_namemydata如果不设置默认会取 pipeline 名并追加_dataset后缀。需要特别区分两个概念数据库 schema本文语境指目的端数据库中数据的结构与组织方式包括表定义和表间关系dlt Schema特指 dlt pipeline 内部规范化数据的格式与结构表结构、列定义的元数据对象两者同名但含义不同阅读文档时不要混淆。Resource 到表的映射pipeline 定义中的每个 resource 都会在目的端对应一张表。上例中只有一个users因此得到一张表mydata.users其中mydata是 schema 名users是表名。table_name是显式设置的若不设置表名默认取 resource 名。等价写法 dlt.resource def users(): yield [ {id: 1, name: Alice}, {id: 2, name: Bob} ] pipeline dlt.pipeline( pipeline_namequick_start, destinationduckdb, dataset_namemydata ) load_info pipeline.run(users)结果与上面完全相同——不需要向pipeline.run显式传table_nameusers表会隐式地以dlt.resource装饰的 resource 函数名users()命名。特殊说明dlt 还会创建若干跟踪 pipeline 状态的内部表它们以_dlt_为前缀不会出现在dlt pipeline show命令的输出中但直连数据库时可以查到。它们的完整结构见下文 dlt 的内部管理表 一节。嵌套数据根表与嵌套表更复杂的例子数据中包含 Python 列表嵌套的对象。import dlt data [ { id: 1, name: Alice, pets: [ {id: 1, name: Fluffy, type: cat}, {id: 2, name: Spot, type: dog} ] }, { id: 2, name: Bob, pets: [ {id: 3, name: Fido, type: dog} ] } ] pipeline dlt.pipeline( pipeline_namequick_start, destinationduckdb, dataset_namemydata ) load_info pipeline.run(data, table_nameusers)运行后会在目的端创建两张表users根表 root table和users__pets嵌套表 nested table。users存放顶层数据users__pets存放来自 Python 列表的嵌套数据。表内容可能如下mydata.usersidname_dlt_id_dlt_load_id1AlicewX3f5vn801W16A1234562350.984172BobrX8ybgTeEmAmmA1234562350.98417mydata.users__petsidnametype_dlt_id_dlt_parent_id_dlt_list_idx1Fluffycatw1n0PEDzuP3grwwX3f5vn801W16A02Spotdog9uxh36VU9lqKpwwX3f5vn801W16A13Fidodogpe3FVtCWz8VuNArX8ybgTeEmAmmA0dlt 在推断数据库 schema 时会把 Python 对象结构例如解析后的 JSON 文件映射为嵌套表并在表之间建立引用。具体规则所有根表和嵌套表每一行都包含一个名为_dlt_id的唯一列row key行主键每张嵌套表都包含_dlt_parent_id列引用父表中特定行的_dlt_idparent key来自 Python 列表的行其在列表中的位置由_dlt_list_idx保存对以mergewrite disposition 加载的嵌套表还会增加root key列_dlt_root_id把子表行引用回根表的对应行。这些引用在源码中是有正式定义的dlt/common/schema/typing.py#L46-L60 中定义了C_DLT_ID _dlt_id、C_DLT_LOAD_ID _dlt_load_id以及_dlt_parent子表→父表隐式引用、_dlt_root后代表→根表隐式引用、_dlt_load根表→_dlt_loads隐式引用三个引用标签dlt/common/schema/utils.py#L1416-L1437 中dlt_id_column()给出_dlt_id的列定义text类型、precision: 64、非空、唯一、row_key: Truedlt_load_id_column()给出_dlt_load_id的列定义text、非空。更详细的嵌套引用、row key 与 parent key 机制见 schema 文档。命名约定表名与列名pipeline 运行期间dlt 会对表名和列名做归一化详见 命名约定文档确保其兼容目的端数据库可接受的格式。来自源数据的所有名称都会转换为 snake_case且只包含字母和数字。注意目的端中的名称可能与原始输入略有差异。从源码结构看默认的归一化规则实现在 dlt/common/normalizers/naming/snake_case.pydlt 默认使用的命名约定其类文档字符串列出的具体规则包括去除首尾空格删除除 ASCII 字母数字和下划线外的所有字符替换为下划线名称以数字开头时前置_连续的多个下划线合并为一个末尾的下划线替换为x与*替换为x、-替换为_、替换为a、|替换为l。该文件还明确说明使用__双下划线作为表之间父子关系和扁平化列名的分隔符——这正是嵌套表users__pets中双下划线的来源。此外 dlt 还支持多种可替换的命名约定实现如 sql_cs_v1.py、sql_ci_v1.py、duck_case.py 等位于 dlt/common/normalizers/naming/ 目录。Variant 列处理类型不一致的数据当同一字段的数据类型不一致时dlt 会把数据分发到多个variant 列。例如某 resource比如 JSON 文件有个字段answer第一次加载时只有布尔值则目的端得到BOOLEAN类型的answer列若下一次加载出现整数和字符串值不一致的数据会分别进入answer__v_bigint和answer__v_text列。variant 列的通用命名规则是original name__v_typeoriginal_name为发生类型冲突的既有列名type为存储在 variant 中的数据类型的名字。列属性中的variant标记在 dlt/common/schema/typing.py#L66-L80 的TColumnProp定义中可见。Load package 与 load ID每次 pipeline 执行会生成一个或多个 load package一个 package 通常包含该次运行中来自 source 所有 resource 的数据。每个 package 由唯一的load_id标识。该load_id会被写入两类位置顶层数据表中的_dlt_load_id列见上文示例特殊的_dlt_loads表其中status为 0 表示加载过程已完全完成。继续向同一目的端加载新数据data [ { id: 3, name: Charlie, pets: [] }, ]pipeline 其余定义不变。这次运行会创建带有新load_id的新 load package并把数据追加到既有表中。users表变成mydata.usersidname_dlt_id_dlt_load_id1AlicewX3f5vn801W16A1234562350.984172BobrX8ybgTeEmAmmA1234562350.984173Charlieh8lehZEvT3fASQ1234563456.12345_dlt_loads表则变为mydata._dlt_loadsload_idschema_namestatusinserted_atschema_version_hash1234562350.98417quick_start02023-09-12 16:45:51.1786500aOEb...Qekd/581234563456.12345quick_start02023-09-12 16:46:03.1066200aOEb...Qekd/58_dlt_loads表追踪已完成的加载并支持在其上串联转换。许多目的端不支持分布式长事务例如 Amazon Redshift此时用户可能看到部分加载的数据。可以把它过滤掉任何load_id不在_dlt_loads中的行都尚未完成加载。同样的方法也可以用来识别并删除永远未完成 package 的数据。其他相关实践对每次加载你可以检测异常例如没有数据、某表加载量过大并发送告警见 生产运行文档的 Slack 告警一节上文提到的 dashboard 应用中也提供了一些有用的 load 统计可以利用status列把 转换transformations串联起来第一个转换从status 0的行开始处理完后更新为 1下一个转换从status 1开始并更新为 2每个附加转换依此类推。数据血缘Data lineage数据血缘在 Data Vault 架构大型组织用于跨系统表示同一业务流程、对数据血缘有强需求的数仓模式或问题排查场景下尤其重要。利用 dlt 开箱即用的 pipeline 名和load_id你可以定位数据的来源和加载时间。你還可以为特定load_id保存完整的血缘信息包括加载的文件列表、错误信息如有、耗时、schema 变更等这对排障很有帮助。Staging 数据集merge 写入的原子性保障前文的示例 pipeline 一直使用appendwrite disposition——每次运行都把数据追加到既有表中。当改用 merge write disposition时dlt 会创建一个 staging 数据库 schema 用于暂存数据默认命名为dataset_name_staging详见 staging 文档其包含与目的端 schema 相同的表集合。运行 pipeline 时staging 表中的数据会在单个原子事务中被写入目的端表。把 pipeline 改为mergeimport dlt dlt.resource(primary_keyid, write_dispositionmerge) def users(): yield [ {id: 1, name: Alice 2}, {id: 2, name: Bob 2} ] pipeline dlt.pipeline( pipeline_namequick_start, destinationduckdb, dataset_namemydata ) load_info pipeline.run(users)运行后目的端会出现名为mydata_staging的 schema。检查其中的表会发现mydata_staging.users与上一节mydata.users相同。源码佐证staging 数据集名的默认布局定义在 dlt/common/destination/client.py#L316staging_dataset_name_layout: str %s_staging同一文件中normalize_staging_dataset_name()、with_staging_dataset()等方法负责在加载流程中切换 staging 目的端并执行回写。运行后表内容可能如下mydata_staging.usersidname_dlt_id_dlt_load_id1Alice 2wX3f5vn801W16A2345672350.984172Bob 2rX8ybgTeEmAmmA2345672350.98417mydata.usersidname_dlt_id_dlt_load_id1Alice 2wX3f5vn801W16A2345672350.984172Bob 2rX8ybgTeEmAmmA2345672350.984173Charlieh8lehZEvT3fASQ1234563456.12345可以看到mydata.users同时包含了之前 pipeline 运行的数据和本次 merge 后的数据id 1、2 被更新为Alice 2、Bob 2。Dev mode版本化versioned数据集在dlt.pipeline调用中把dev_mode参数设为True时dlt 会创建版本化数据集每次运行 pipeline数据都会加载到一个新的数据集新的数据库 schema中数据集名是你提供的dataset_name加一个基于日期时间的后缀。import dlt data [ {id: 1, name: Alice}, {id: 2, name: Bob} ] pipeline dlt.pipeline( pipeline_namequick_start, destinationduckdb, dataset_namemydata, dev_modeTrue # -- add this line ) load_info pipeline.run(data, table_nameusers)每次运行都会在目的端数据库创建一个带日期时间后缀的新 schema第一次运行可能是mydata_20230912064403第二次是mydata_20230912064407依此类推。数据被加载到这些新 schema 的表中。源码佐证dev_mode参数定义于 dlt/pipeline/init.py 的dlt.pipeline()签名默认False文档字符串说明其语义为“每个同pipeline_name的 pipeline 实例运行时从零开始并把数据加载到独立的数据集”旧参数full_refresh已标记弃用、等价于dev_mode。Pipeline类在 dlt/pipeline/pipeline.py 中把dev_mode持久化到状态dev_mode字段并在 attach 已有状态时恢复该标志。dlt 的内部管理表dlt 会自动在目的端 schema 中创建内部管理表用于跟踪 pipeline 运行、支持增量加载、管理 schema 版本均以_dlt_为前缀。表名常量集中定义在 dlt/common/schema/typing.py#L39-L43VERSION_TABLE_NAME _dlt_version、LOADS_TABLE_NAME _dlt_loads、PIPELINE_STATE_TABLE_NAME _dlt_pipeline_state、DLT_NAME_PREFIX _dlt。_dlt_loads加载历史追踪每次 pipeline 执行都会向该表插入一行带唯一load_id。它记录哪些 load 已完成并支持串联转换。Column nameTypeDescriptionload_idSTRING加载任务的唯一标识schema_nameSTRING加载时使用的 schema 名schema_version_hashSTRINGschema 版本哈希statusINTEGER加载状态0表示完成inserted_atTIMESTAMP该 load 被记录的时间只有status 0的行才算完成其他值代表未完成或被打断的加载。status 列还可以用于协调多步转换。列定义与源码一致dlt/common/schema/utils.py#L1369-L1413 的loads_table()中load_id为textprecision 64非空、schema_name为可空text、status为bigint非空描述 0 success、inserted_at为timestamp非空、schema_version_hash为可空text且该表write_disposition skip即 dlt 跳过对它的常规写入逻辑仅在加载完成时追加记录。_dlt_pipeline_statepipeline 状态与检查点该表保存 pipeline 每次运行的内部状态使增量加载成为可能并在上次运行被打断时从断点恢复。Column nameTypeDescriptionversionINTEGER该状态条目的版本engine_versionINTEGER使用的 dlt 引擎版本pipeline_nameSTRINGpipeline 名称stateSTRING or BLOB序列化后的 pipeline 状态 Python 字典created_atTIMESTAMP状态条目创建时间version_hashSTRING用于检测状态变更的哈希_dlt_load_idSTRING引用_dlt_loads中的关联 load_dlt_idSTRINGpipeline 状态行的唯一标识state列包含的序列化 Python 字典包括增量进度如最后处理的条目或时间戳、转换检查点、source 相关的元数据与设置。这使 dlt 能恢复被中断的 pipeline、避免重复加载已处理数据保证 pipeline 幂等且高效。version_hash在每次更新时重新计算——dlt 正是依赖该表实现 last-value 增量加载即使某次运行失败或中断下次运行也会从正确的检查点继续。源码佐证dlt/common/schema/utils.py#L1440-L1494 的pipeline_state_table()给出了列定义state列描述为 Compressed JSON representation of the pipeline state即状态以压缩 JSON 存储表默认write_disposition append_dlt_id列仅在add_dlt_idTrue时追加。_dlt_versionschema 版本追踪该表记录 pipeline 使用过的所有 schema 版本历史。每当 dlt 更新 schema例如新增列或表时都会向该表写入一条新记录。Column nameTypeDescriptionversionINTEGERschema 的数字版本号engine_versionINTEGER使用的 dlt 引擎版本inserted_atTIMESTAMPschema 版本条目创建时间schema_nameSTRINGschema 名称version_hashSTRING表示 schema 内容的唯一哈希schemaSTRING or JSONJSON 格式的完整 schema保留历史 schema 定义保证了旧数据仍可按原定义读取新数据使用更新后的 schema 规则向后兼容性得以维持。该表同时支持排障与兼容性检查——可以追踪任意一次 load 使用的是哪个 schema 版本和引擎版本帮助调试并确保数据模型安全演进。_dlt_loads.schema_version_hash与_dlt_version.version_hash之间、schema_name与schema_name之间的关联在 dlt/common/schema/utils.py#L1254-L1297 中以显式表引用_dlt_schema_version、_dlt_schema_name的形式建立。向非 dlt 创建的既有表加载数据你也能把 dlt 数据加载到目的端数据集中已存在、并非由 dlt 创建的表但需注意三种情形表存在但无数据多数情况下加载会顺利成功——dlt 会创建所需列并插入数据。dlt 只认识其内部 schema 中发现或提供的列目的端上 dlt 未知的列会保留在表中但对 dlt 不可见通常没有问题。表存在且列名与 dlt 发现的列同名但数据类型不匹配加载会失败。你必须先在目的端修改该列或把入站数据中的列名改成别的名字以避免冲突。表存在且已有数据加载最初可能失败因为 dlt 会创建包含必需元数据的non-nullable列。很多数据库不允许在已有数据的表上创建non-nullable列现有行的初始值无法推断。你需要手动在既有表上以正确类型创建这些列、设为nullable再为现有行填充值。部分数据库允许在同一命令中创建non-nullable列并为现有行取默认值。需要创建的列为nametype_dlt_load_idtext/string/varchar_dlt_idtext/string/varchar对于嵌套表可能还需要创建nametype_dlt_parent_idtext/string/varchar_dlt_root_idtext/string/varchar小结dlt 在目的端数据库中构建了一套自描述、可追溯的表组织体系以dataset_name命名的 schema 承载各 resource 对应的表嵌套 Python 结构被拆分为由_dlt_id/_dlt_parent_id/_dlt_list_idx/_dlt_root_id串联的根表与嵌套表snake_case命名约定__作为父子分隔符保证标识符兼容性load_id贯穿数据表与_dlt_loads支撑血缘追踪、异常过滤与转换串联merge 场景下的dataset_name_staging数据集提供原子写入dev_mode则让每次运行产出独立版本化数据集。理解这套结构后你可以直接对目的端数据库写 SQL 完成审计、排障与二次开发而不必依赖 dlt 的 Python 接口。【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考