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_name='quick_start', destination='duckdb', dataset_name='mydata' ) load_info = pipeline.run(data, table_name="users")运行后,dlt 会在目的端数据库(这里是 DuckDB 内存数据库)中创建一个数据库 schema,以及其中一张名为users的表,并把你 source 中的数据写入其中。其他数据库目的端的行为与概念基本一致。
小技巧:可以使用 dlt pipeline CLI 的
show命令查看目的端数据库中的表,例如dlt pipeline show;dashboard 文档中也有相关说明。
数据库 schema 的命名
数据库 schema 是承载你已加载数据的一组表。schema 名与 pipeline 定义中提供的dataset_name相同:本例中显式设置了dataset_name='mydata';如果不设置,默认会取 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_name='quick_start', destination='duckdb', dataset_name='mydata' ) load_info = pipeline.run(users)结果与上面完全相同——不需要向pipeline.run显式传table_name="users",表会隐式地以@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_name='quick_start', destination='duckdb', dataset_name='mydata' ) load_info = pipeline.run(data, table_name="users")运行后会在目的端创建两张表:users(根表 root table)和users__pets(嵌套表 nested table)。users存放顶层数据,users__pets存放来自 Python 列表的嵌套数据。表内容可能如下:
mydata.users
| id | name | _dlt_id | _dlt_load_id |
|---|---|---|---|
| 1 | Alice | wX3f5vn801W16A | 1234562350.98417 |
| 2 | Bob | rX8ybgTeEmAmmA | 1234562350.98417 |
mydata.users__pets
| id | name | type | _dlt_id | _dlt_parent_id | _dlt_list_idx |
|---|---|---|---|---|---|
| 1 | Fluffy | cat | w1n0PEDzuP3grw | wX3f5vn801W16A | 0 |
| 2 | Spot | dog | 9uxh36VU9lqKpw | wX3f5vn801W16A | 1 |
| 3 | Fido | dog | pe3FVtCWz8VuNA | rX8ybgTeEmAmmA | 0 |
dlt 在推断数据库 schema 时,会把 Python 对象结构(例如解析后的 JSON 文件)映射为嵌套表,并在表之间建立引用。具体规则:
- 所有(根表和嵌套表)每一行都包含一个名为
_dlt_id的唯一列(row key,行主键); - 每张嵌套表都包含
_dlt_parent_id列,引用父表中特定行的_dlt_id(parent 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: True),dlt_load_id_column()给出_dlt_load_id的列定义(text、非空)。更详细的嵌套引用、row key 与 parent key 机制见 schema 文档。
命名约定:表名与列名
pipeline 运行期间,dlt 会对表名和列名做归一化(详见 命名约定文档),确保其兼容目的端数据库可接受的格式。来自源数据的所有名称都会转换为 snake_case,且只包含字母和数字。注意:目的端中的名称可能与原始输入略有差异。
从源码结构看,默认的归一化规则实现在 dlt/common/normalizers/naming/snake_case.py(dlt 默认使用的命名约定),其类文档字符串列出的具体规则包括:
- 去除首尾空格;
- 删除除 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_<type>:original_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.users
| id | name | _dlt_id | _dlt_load_id |
|---|---|---|---|
| 1 | Alice | wX3f5vn801W16A | 1234562350.98417 |
| 2 | Bob | rX8ybgTeEmAmmA | 1234562350.98417 |
| 3 | Charlie | h8lehZEvT3fASQ | 1234563456.12345 |
_dlt_loads表则变为:
mydata._dlt_loads
| load_id | schema_name | status | inserted_at | schema_version_hash |
|---|---|---|---|---|
| 1234562350.98417 | quick_start | 0 | 2023-09-12 16:45:51.17865+00 | aOEb...Qekd/58= |
| 1234563456.12345 | quick_start | 0 | 2023-09-12 16:46:03.10662+00 | aOEb...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 改为merge:
import dlt @dlt.resource(primary_key="id", write_disposition="merge") def users(): yield [ {'id': 1, 'name': 'Alice 2'}, {'id': 2, 'name': 'Bob 2'} ] pipeline = dlt.pipeline( pipeline_name='quick_start', destination='duckdb', dataset_name='mydata' ) load_info = pipeline.run(users)运行后,目的端会出现名为mydata_staging的 schema。检查其中的表,会发现mydata_staging.users与上一节mydata.users相同。
源码佐证:staging 数据集名的默认布局定义在 dlt/common/destination/client.py#L316,
staging_dataset_name_layout: str = "%s_staging";同一文件中normalize_staging_dataset_name()、with_staging_dataset()等方法负责在加载流程中切换 staging 目的端并执行回写。
运行后表内容可能如下:
mydata_staging.users
| id | name | _dlt_id | _dlt_load_id |
|---|---|---|---|
| 1 | Alice 2 | wX3f5vn801W16A | 2345672350.98417 |
| 2 | Bob 2 | rX8ybgTeEmAmmA | 2345672350.98417 |
mydata.users
| id | name | _dlt_id | _dlt_load_id |
|---|---|---|---|
| 1 | Alice 2 | wX3f5vn801W16A | 2345672350.98417 |
| 2 | Bob 2 | rX8ybgTeEmAmmA | 2345672350.98417 |
| 3 | Charlie | h8lehZEvT3fASQ | 1234563456.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_name='quick_start', destination='duckdb', dataset_name='mydata', dev_mode=True # <-- add this line ) load_info = pipeline.run(data, table_name="users")每次运行都会在目的端数据库创建一个带日期时间后缀的新 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-L43:VERSION_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 name | Type | Description |
|---|---|---|
load_id | STRING | 加载任务的唯一标识 |
schema_name | STRING | 加载时使用的 schema 名 |
schema_version_hash | STRING | schema 版本哈希 |
status | INTEGER | 加载状态,0表示完成 |
inserted_at | TIMESTAMP | 该 load 被记录的时间 |
只有status = 0的行才算完成;其他值代表未完成或被打断的加载。status 列还可以用于协调多步转换。
列定义与源码一致:dlt/common/schema/utils.py#L1369-L1413 的loads_table()中,load_id为text(precision 64,非空)、schema_name为可空text、status为bigint(非空,描述 "0 = success")、inserted_at为timestamp(非空)、schema_version_hash为可空text,且该表write_disposition = "skip"(即 dlt 跳过对它的常规写入逻辑,仅在加载完成时追加记录)。
_dlt_pipeline_state:pipeline 状态与检查点
该表保存 pipeline 每次运行的内部状态,使增量加载成为可能,并在上次运行被打断时从断点恢复。
| Column name | Type | Description |
|---|---|---|
version | INTEGER | 该状态条目的版本 |
engine_version | INTEGER | 使用的 dlt 引擎版本 |
pipeline_name | STRING | pipeline 名称 |
state | STRING or BLOB | 序列化后的 pipeline 状态 Python 字典 |
created_at | TIMESTAMP | 状态条目创建时间 |
version_hash | STRING | 用于检测状态变更的哈希 |
_dlt_load_id | STRING | 引用_dlt_loads中的关联 load |
_dlt_id | STRING | pipeline 状态行的唯一标识 |
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_id=True时追加。
_dlt_version:schema 版本追踪
该表记录 pipeline 使用过的所有 schema 版本历史。每当 dlt 更新 schema(例如新增列或表)时,都会向该表写入一条新记录。
| Column name | Type | Description |
|---|---|---|
version | INTEGER | schema 的数字版本号 |
engine_version | INTEGER | 使用的 dlt 引擎版本 |
inserted_at | TIMESTAMP | schema 版本条目创建时间 |
schema_name | STRING | schema 名称 |
version_hash | STRING | 表示 schema 内容的唯一哈希 |
schema | STRING or JSON | JSON 格式的完整 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列并为现有行取默认值。需要创建的列为:
| name | type |
|---|---|
| _dlt_load_id | text/string/varchar |
| _dlt_id | text/string/varchar |
对于嵌套表,可能还需要创建:
| name | type |
|---|---|
| _dlt_parent_id | text/string/varchar |
| _dlt_root_id | text/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),仅供参考