恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
dlt 技术指南:用 @dlt.hub.transformation 构建 Eager/Lazy 双模数据转换
首页
资讯中心
/
dlt 技术指南:用 @dlt.hub.transformation 构建 Eager/Lazy 双模数据转换
dlt 技术指南:用 @dlt.hub.transformation 构建 Eager/Lazy 双模数据转换
发布时间:2026/9/17 20:40:21
dlt 技术指南用 dlt.hub.transformation 构建 Eager/Lazy 双模数据转换【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt本篇基于 dlt 仓库中 dltHub Transformations 文档 展开讲解如何用一个 Python 装饰器dlt.hub.transformation定义数据转换并让同一份代码既能在本地计算DuckDB、Pandas、Polars、Arrow上 eager 执行也能作为下推 SQL 在数仓中 lazy 执行。读完后你可以掌握增量游标式转换的完整写法、跨数据集 Join 的执行位置选择机制、SQL 方言自动编译、Schema 演进与列级血缘转发以及 dlt 处理这类转换的 extract → normalize → load 全生命周期。一、什么是 dltHub TransformationsdltHub 在dlt之上扩展了dlt.hub.transformation装饰器。它与dlt.resource的区别在于转换函数是一个接收dlt.Dataset的 generatoryield 出的不是数据行而是Ibis 表达式、Relation对象或一段 SQL 查询。dlt 会将该查询编译到目标数据源的 SQL 方言并把结果物化到目标表中。核心设计可以概括为一句话同一个转换函数由 dlt 根据源与目标是否位于同一物理位置自动决定数据在哪里执行。当源和目标共享同一物理位置例如同一个 DuckDB 文件、同一个 MotherDuck 账号时数据不需要传输到运行 pipeline 的机器反之dlt 会把查询结果提取为 Parquet/Arrow 数据在本地机器上执行再走一次常规的 dlt load。1.1 插件机制dlt.hub 是如何加载的从源码看dlt.hub是一个可选插件入口其加载逻辑在 dlt/hub/init.py 中try: from dlthub import transformation, runner from . import current __found__ True __all__ (transformation, current, runner) except ImportError as import_exc: __exception__ import_excdlt.hub在模块导入时尝试加载商业插件包dlthub中的transformation与runner如果未安装模块本身不会报错而是在你访问dlt.hub.transformation等属性时通过模块级__getattr__抛出一个带有明确安装指引的MissingDependencyException提示执行pip install dlt[hub]。此外dlt/hub/current.py 直接 re-exportdlthub.currentdlt/hub/run.py 则复用 workspace 的部署装饰器job、pipeline_run、trigger等供 dltHub Platform 的部署场景使用。测试用例印证了这一契约tests/hub/test_plugin_import.py 断言插件存在时dlt.hub.__found__为真、访问不存在的属性抛普通AttributeErrortests/hub/test_transformations.py 则验证了dlt.hub.transformation装饰器返回的对象携带正确的资源名transformation.name get_even_rows且dlt.dataset()可以构造 mock 数据集用于测试。适用前提dlt.hub系列功能是 dltHub 商业特性其使用受商业 dltHub License 约束安装对应版本的dlthub插件包后方可使用。二、最小可运行示例增量追加 recent_orders原文档给出的核心示例是把orders表中新增的行增量追加到recent_orders每次运行从上一次处理到的created_at之后继续。完整代码如下import dlt import pendulum dlt.hub.transformation( write_dispositionmerge, primary_keyid, incrementaldlt.sources.incremental( created_at, initial_valuependulum.datetime(2000, 1, 1, tzUTC), ), ) def recent_orders(dataset: dlt.Dataset) - Any: yield dataset.table(orders) pipeline.run(recent_orders(pipeline.dataset()))逐点解析装饰器参数与dlt.resource一致write_dispositionappend/replace/merge等、primary_key、table_name、columns提示、incremental均可直接使用转换与常规资源享受完全相同的 load 语义。隐式游标注意函数体内 yield 的是“裸”关系dataset.table(orders)并没有手动加过滤条件。incremental声明在装饰器上时dlt 会自动把游标过滤应用到 yield 出的查询上。这是隐式游标形式另一种显式形式是把incremental作为函数参数传入并在体内调用.incremental(window)见第五节。merge primary_key 组合write_dispositionmerge配合primary_keyid使边界游标值对应的行即使被重复读到也会被合并去重适合需要“边界行立即可见”的场景。一个更朴素的入门示例完整 quickstart 见 dltHub transformations 文档dlt.hub.transformation def copied_customers(dataset: dlt.Dataset) - Any: customers_table dataset[customers] yield customers_table.order_by(name).limit(5) fruitshop_pipeline.run(copied_customers(fruitshop_pipeline.dataset()))dlt 会检测到转换读写的是同一个 DuckDB 数据集于是把整个转换作为 SQL 在目标端执行——没有任何数据进出运行 pipeline 的机器同时自动为新表copied_customers执行 schema 演进。三、Eager 还是 Lazy执行位置的自动判定这是dlt.hub.transformation最有价值的特性执行位置不是写死的而是由 dlt 根据输入输出数据所在位置的拓扑关系推断。3.1 三种典型情形情形执行方式数据是否离开目标端输入与输出在同一数据位置同一 DuckDB 文件、同一 MotherDuck 账号、同一 bucket目标端纯 SQL 执行model job否输入与输出在不同数据位置但输出引擎可读输入如filesystem数据集 attach 进 DuckDB 输出仓内 SQL ATTACH语句否不同引擎如本地 DuckDB → Postgres或你主动 yield DataFrame/Arrow 表eager 物化本地执行查询把结果作为 Parquet/Arrow 走常规 dlt load是跨引擎示例来自 hub 文档把同一个copied_customers转换交给指向 Postgres 的 pipeline 运行# different engine (DuckDB → Postgres) duck_p dlt.pipeline(fruitshop_warehouse, destinationpostgres) duck_p.run(copied_customers(fruitshop_pipeline.dataset()))dlt 会用查询从 DuckDB 提取出数据Parquet 文件再对 Postgres 执行一次常规 dlt load。转换函数代码一个字都不用改——你只需要按场景选择 pipeline 指向哪个 destination。3.2 多数据集 Join 的执行位置选择转换函数可以接收多个dlt.Dataset作为参数dlt 会检查输入与输出分别在哪里然后选择执行方式# crm 与 sales 是同一 DuckDB 文件里的两个数据集marts 是输出 dlt.hub.transformation(table_nameuser_orders) def user_orders(crm: dlt.Dataset, sales: dlt.Dataset) - Any: yield crm[users].join(sales[orders], onusers.id orders.user_id) marts_pipeline.run(user_orders(crm_pipeline.dataset(), sales_pipeline.dataset()))若输出引擎既能读又能写所有输入一个 DuckDB 库或一个 MotherDuck 账号Join 作为model job 在仓内执行否则走eager 物化在本地机器上跑完查询、以数据形式加载结果。对 DuckDB 系引擎duckdb、ducklake、motherduck、lance/lancedb以及支持file/s3/az/abfss/hf协议的filesystem含 delta/iceberg 开放表格式dlt 还能把输入数据集以attach alias的方式挂到输出端的 DuckDB 引擎上执行 Join。两个值得注意的细节MotherDuck 的输入 attach 在本地完成除另一个 MotherDuck 库之外的输入都在你的本地会话 attach输入端的凭据留在本地、不会被上传到 MotherDuck只有数据行随查询流动强制 eager只要 yield 的是物化结果Arrow 表、DataFrame而不是 Relationdlt 就不会创建 model job也不会序列化任何凭据dlt.hub.transformation(table_nameuser_orders_eager) def user_orders_eager(warehouse: dlt.Dataset, orders: dlt.Dataset) - Any: joined warehouse[users].join(orders[orders], onusers.id orders.user_id) yield joined.arrow() # 本地执行无 model job、无凭据序列化跨位置 Join 若需要凭据MotherDuck token、云桶密钥等dlt 会把这些语句加密后写入 model job 的.model文件密钥来自 pipeline 的加密种子若未显式配置pipeline_salt进程重启后解密会失败此时按提示在secrets.toml中设置pipelines.pipeline_name.pipeline_salt使密钥可复现。四、三种写法Ibis 表达式、原生 SQL、DataFrame/Arrow转换函数 yield 的内容决定 dlt 如何处理它第一个 yield 项是合法的 SQL 查询或 Relation 对象时dlt 按 SQL 转换处理否则如 DataFrame、Arrow 表、普通 Python 对象退化为常规 resource 语义。4.1 Ibis 表达式dlt.hub.transformation(nameorders_per_user, write_dispositionmerge) def orders_per_user(dataset: dlt.Dataset) - Any: purchases dataset.table(purchases).to_ibis() yield purchases.group_by(purchases.customer_id).aggregate( order_countpurchases.id.count() )dlt.Dataset上的 Ibis 表达式.table()、.to_ibis()、join、group_by、aggregate、order_by、limit 等是写转换的主力详见 dataset 文档。4.2 原生 SQL 与跨方言编译dlt.hub.transformation def enriched_purchases(dataset: dlt.Dataset) - Any: enriched dataset( SELECT customers.name, purchases.quantity FROM purchases JOIN customers ON purchases.customer_id customers.id ) yield enriched通过dataset(SQL)可以直接写原始 SQL其中标识符是dlt schema 中的表名/列名而非目标数据库的物理 schema 名。若你的查询是用另一种方言写的传入query_dialect参数dlt 会解析该方言并重新生成目标方言的 SQL——例如 duckdb 的||拼接和LIMIT会被编译成 mssql 的与TOP标识符自动套用目标端的引号规则。这意味着一份为某数仓写好的查询可以直接在另一个数仓上运行。4.3 Pandas / Polars / Arrowdlt.hub.transformation def enriched_purchases(dataset: dlt.Dataset) - Any: purchases dataset.table(purchases).df() customers dataset.table(customers).df() result purchases.merge(customers, left_oncustomer_id, right_onid) yield result[[name, quantity]]yield DataFrame 或 Arrow 表后该转换资源按常规 resource 处理dlt 不转发列级 hintsDataFrames 被视为普通数据。此行为文档注明可能在未来版本变化。这条路径天然就是 eager 执行——计算发生在运行 pipeline 的机器上适合在数据进入昂贵数仓之前做预聚合、预过滤来节省仓内算力文档中给出了rest_api源 → 本地 DuckDB 中转 → 仅把聚合结果写入 Postgres 的完整示例。五、增量转换游标与调度窗口当源数据持续增长每次全量重跑转换既慢又贵。增量转换让每次运行只处理源数据的正确切片切片由**游标cursor**定义一个列其值与某个范围比较以决定行是否在本次范围内。常见选择是created_at、updated_at、递增的id或 dlt 管理的_dlt_loads.inserted_at。dlt 提供两种选片方式5.1 调度器窗口scheduler interval当编排器决定每次运行的时间范围时使用cron 调度、重试、分区回填的自然选择。把游标声明为函数参数并设置allow_external_schedulersTrueimport dlt from dlt.common.pendulum import pendulum dlt.hub.transformation(write_dispositionreplace) def orders_window( dataset: dlt.Dataset, window: dlt.sources.incremental[pendulum.DateTime] dlt.sources.incremental( created_at, initial_valuependulum.datetime(2000, 1, 1, tzUTC), allow_external_schedulersTrue, range_startclosed, range_endopen, ), ) - Any: yield dataset.table(orders).incremental(window)在 dltHub Platform 上调度器设置DLT_INTERVAL_START与DLT_INTERVAL_END环境变量dlt 用它们过滤源数据平台侧机制见 triggers 文档。窗口是[start, end)半开区间给定 2026-01-01 至 2026-01-10 每天一行的orders表窗口[2026-01-05, 2026-01-10)会写入 id 5~9id 10 因开区间端被排除。重放同一个窗口产生完全相同的输入因此该模式天然幂等非常适合分区回填和重试。5.2 有状态游标从上次运行继续每次运行需接着上次成功位置继续时使用有状态游标——不需要外部调度器dlt 内部持久化游标状态dlt.hub.transformation( write_dispositionappend, primary_keyid, incrementaldlt.sources.incremental( created_at, initial_valuependulum.datetime(2000, 1, 1, tzUTC), range_startopen, ), ) def recent_orders(dataset: dlt.Dataset) - Any: yield dataset.table(orders)假设orders分两批加载id 1..3 对应 01-01~01-03id 4..5 对应 01-04~01-05首次运行没有last_value从initial_value开始写入 3 行并把last_value推进到2026-01-03下一次运行看到 2 行落在last_value之后追加它们并把last_value推进到2026-01-05。范围range语义决定边界行行为——每次有状态运行在抽取时计算MAX(cursor)并以此为范围端点因此边界行MAX 所在行的处理方式由 range 设置决定range_startopenappend 推荐边界行在记录它的这次运行中加载后续运行不再读它默认range_startclosedrange_endopen边界行被推迟下一次观察到更大游标值的运行时恰好加载一次——append 不会产生重复但最新行要等一个周期range_startclosedrange_endclosedwrite_dispositionmerge 主键边界行立即加载且每次运行重读由 merge 去重——适合“共享边界游标值的迟到行必须立刻可见”的场景若primary_key恰好等于游标列声明游标值唯一边界行总是立即加载且不重放。5.3 游标列的选择与安全规则只追加append-only数据用created_at或递增id可变数据用随每次变更而变化的updated_at源表没有业务时间戳时用_dlt_loads.inserted_at按加载时间处理。点分游标路径让 dlt 从基础表 join 到_dlt_loads并过滤join 仅用于过滤不会把_dlt_loads列加进目标表状态安全规则有状态 relation 增量上出现LIMIT会被拒绝从受限结果推进状态可能跳过行SQL 游标支持max/min两种 last-value 函数自定义 Pythonlast_value_func无法下推 SQL空值处理遵循on_cursor_value_missinginclude追加OR cursor IS NULLexclude追加AND cursor IS NOT NULLraise在查询中途无法抛出、回退为排除空值。内部实现上当转换以 model job 方式执行时dlt 在运行前改写源查询注入游标过滤另外两种情况在抽取阶段过滤源与目标位于不同数据位置或 yield 的是 Python 对象列表、Arrow 表、DataFrame。六、Schema 演进、列级血缘与转换生命周期6.1 先算 schema再执行dlt 在执行转换之前就计算出结果 schema由此可以做到迁移目标 schema新建列/表、在不兼容的 schema 变更如有损的类型变更发生时提前失败、以及把列级 hints 从源表转发到转换后的表。开发期可通过Relation.columns/Relation.columns_schema检视计算出的结果 schema。6.2 列级 hints 转发dlt 会转发特定 hints所有x-annotation...开头的自定义 hint以及nullable、data_type、precision、scale、timezone类型 hints。例如源表name列标记了x-annotation-pii: True经 Join 生成的新表中name列会自动保留 PII 标记——PII 审计可以在转换层之后追溯。两个边界primary_key、merge_keys这类 hints 不会转发需要用装饰器的columns参数显式声明因为 dlt 不知道你如何消费转换结果表由多个来源列组合产生的列如拼接、SQL 运算结果无法转发 hints。6.3 SQL 转换的 extract → normalize → load 生命周期yieldRelation的转换SQL 转换走一条独立于普通 resource 的生命周期Extractdlt 把 Relation 转成 SQL 字符串连同其源方言保存为.model文件。此阶段的 SQL 就是你的原始查询或Relation.to_sql()的产物尚未添加_dlt_id、_dlt_load_id等 dlt 列。Normalizedlt 读取.model文件并改写查询——_dlt_load_id默认添加若查询中已存在则被替换为当前 load ID 常量可在settings.toml中用[normalize.model_normalizer] add_dlt_load_id false关闭_dlt_id默认不添加可用add_dlt_id true开启生成方式因方言而异如 ClickHouse 用generateUUIDv4()Redshift 用 load ID 与行号的 MD5 哈希SQLite 用lower(hex(randomblob(16)))模拟用 database/dataset 前缀全限定所有标识符、按目标端要求加引号并调整大小写、按命名规范规范化列名、为命名差异生成别名、按目标表 schema 重排列顺序、为目标中存在但查询中没有的列补NULL。Load把归一化后的查询包进INSERT INTO ... SELECT语句在目标端执行。最终效果是形如INSERT INTO my_pipeline_dataset.my_transformation (id, value, _dlt_load_id, _dlt_id) SELECT _dlt_subquery.id, _dlt_subquery.value, 1749134128.17655 AS _dlt_load_id, UUID() AS _dlt_id FROM ( SELECT my_table.id, my_table.value FROM my_pipeline_dataset.my_table AS my_table ) AS _dlt_subquery目标端的 SQL client 执行该语句转换结果直接物化在数据库内。而 Python 型转换DataFrame/Arrow/Polars则走常规 resource 的 extract、normalize、load 流程。七、与 dltHub Platform 的集成要点调度驱动的增量窗口在 dltHub Platform 上cron 触发器为每次运行设置[start, end)区间通过DLT_INTERVAL_START/DLT_INTERVAL_END配合allow_external_schedulersTrue的游标实现幂等重试与分区回填资源语义完整继承write_disposition、primary_key、merge load、以及把多个转换分组进一个dlt.source甚至与普通 resource 混跑在一次 load 中全部可用许可边界dltHub 功能受商业 dltHub License 约束插件包dlthub需以匹配版本安装pip install dlt[hub]。八、验证与延伸阅读仓库内的测试可以印证本文的关键契约tests/hub/test_transformations.py验证dlt.hub.transformation装饰后资源名保留、dlt.dataset(duckdb, mock_dataset)可构造测试用数据集tests/hub/test_plugin_import.py验证插件发现逻辑dlt.hub.__found__、异常透传与dlthubCLI 命令的可用性。延伸阅读均为仓库内文档dltHub Transformations 完整指南quickstart、跨数据集 Join、DuckDB attach 细节、in-transit 转换示例dltHub Transformations 生态入口文档本文的主体来源dltHub Platform triggers调度器窗口与触发器机制dltHub License许可范围与再分发条款。总结dlt.hub.transformation的本质是“把已加载的数据当作新的数据源来加工”——你用与普通 resource 相同的 API 心智写转换dlt 负责编译方言、计算 schema、选择 eager/lazy 执行路径、管理增量游标状态并转发列级血缘让“加载后的数据再加工”这一环节无需离开熟悉的 dlt pipeline 体系。【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考