使用 dltHub 托管编排部署 dlt 数据管道:从脚手架、部署到定时触发与刷新级联
【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy 🛠️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt
dltHub 是 dlt(data load tool)项目的托管编排平台,围绕@dlt.hub.run装饰器构建了一套"代码即编排"的托管调度体系:调度计划和数据依赖关系直接写在 Python 代码里,而不是保存在独立的 DAG 文件或 YAML 中。本文以仓库文档 orchestrate-with-dlthub.md 为主线,带你走完"安装脚手架 → ad-hoc 运行 →__deployment__.py声明部署 → 定时触发"四个阶段,并深入触发器全家桶(follow-up 链、freshness 门控、刷新级联、标签批量触发、时区)与底层源码实现,读完即可把一条本地 dlt 管道部署到 dltHub 平台并纳入托管调度。
dltHub 托管编排的核心思想:调度住在代码里
传统编排工具(如 Airflow、Prefect)通常需要为每个任务单独声明 DAG 文件,调度表达式与业务代码分居两处。dltHub 的做法完全不同:触发器和数据依赖都是装饰器参数,随函数定义一起版本化、一起部署。正如 orchestrate-with-dlthub.md 开头所述:
- 调度(cron、间隔、一次性、follow-up)通过
trigger=参数挂到装饰器上; - 数据依赖通过
.success/.fail/.completed触发器属性在作业之间显式建模; - 没有独立的"添加/删除调度" CLI 命令——改装饰器,重新
dlthub deploy,就是修改调度的唯一方式(见 triggers.md)。
底层支撑来自 dlt/hub/run.py:该模块从dlt._workspace.deployment导出job、pipeline、interactive三个装饰器、trigger工厂集合与TJobRunContext类型,是全部托管作业的公共入口。
第一步:安装 CLI 并脚手架化一个 workspace
dltHub 的 CLI 名为dlthub,通过uvx直接拉起最新版本即可,无需手工安装:
uvx dlthub-init@latest如果你还没有uv,需要先按 uv 官方安装指南装好它。该命令会在当前目录生成一个完整的 workspace:
.dlt/.workspace标记文件——它的存在激活workspace 模式(CLI 与平台交互的前提);pyproject.toml——声明依赖,平台运行作业时按此文件安装包;.dlt/配置与 secrets 文件(*.toml)。
关于其他安装途径与初始化路径,可参考 installation.md。CLI 与平台的连接(登录、绑定 workspace)细节见 workspace-setup.md。
第二步:ad-hoc 运行——无需部署文件即可跑通管道
任何常规 dlt 管道都可以不写任何部署文件直接在平台上运行。只要你的pipeline.py里包含一个顶层pipeline.run(...)调用:
uv run dlthub run pipeline.py这条命令的底层行为是:把该文件以 ad-hoc 方式部署到平台、立即执行,并把平台上的运行日志流式回传到你的终端。它非常适合一次性运行与冒烟测试。
deployments.md 给出了同一场景的完整命令族(deployments.md):
# 先在本地跑一遍批处理脚本,提前暴露缺依赖或配置错误 dlthub local run fruitshop_pipeline.py # 部署并在云端运行(使用 prod profile) dlthub run fruitshop_pipeline.py # 流式查看日志直到运行结束 dlthub run fruitshop_pipeline.py -f值得强调的是:ad-hoc 部署不支持定时触发器、follow-up 作业、freshness 约束以及多作业整体部署——这些能力必须使用装饰器 +__deployment__.py的清单式部署。
第三步:添加__deployment__.py,把管道声明为托管作业
要把同一条管道升级为"托管作业"(managed job),需要两步:给入口函数加上装饰器,并在 workspace 的__deployment__.py中声明它。
入口文件ingest_breweries.py(使用@run.pipeline装饰器):
import dlt from dlt.hub import run @run.pipeline("ingest_breweries") def ingest_breweries(): pipeline = dlt.pipeline( pipeline_name="ingest_breweries", destination="warehouse", dataset_name="brewery_data", ) pipeline.run(brewery_source())__deployment__.py(workspace 的部署清单):
# __deployment__.py from ingest_breweries import ingest_breweries __all__ = ["ingest_breweries"]然后部署:
uv run dlthub deploy之后按作业名触发运行:
uv run dlthub run ingest_breweries装饰器三兄弟:pipeline / job / interactive
dlt/hub/run.py 对外暴露三个装饰器,用途不同(见 deployments.md):
| 装饰器 | 适用场景 |
|---|---|
@run.pipeline | 绑定到具名dlt.pipeline的批处理作业(获得 pipeline 感知的重试与数据集链接) |
@run.job | 通用批处理作业(数据质量检查、报表、自定义脚本等任意 Python 函数) |
@run.interactive | 长驻 HTTP 服务(notebook、MCP 服务器、Streamlit 应用、REST API) |
__deployment__.py的导入规则(见 deployments.md):
- 函数导入(
from github_pipeline import load_commits):每个被导入且被装饰的函数生成一个作业; - 模块导入(
import github_report_notebook):每个模块生成一个作业,框架自动探测——marimo notebook 变成交互式 notebook 作业、FastMCP 模块变成 MCP 服务器、Streamlit 模块变成仪表盘; __all__:精确列出要部署的名字;没有它清单生成器会扫描__dict__并告警;- 模块 docstring:成为平台仪表盘里的 workspace 描述;
- 也支持在
__deployment__.py内联定义装饰过的作业,适合小型 MCP 服务器或一次性批处理。
装饰器参数全览(源码级)
dlt/_workspace/deployment/decorators.py 中job的完整签名揭示了所有可用参数:
| 参数 | 含义 |
|---|---|
name/section | 作业名 / 配置段,默认取函数名 / 模块名,必须是合法 Python 标识符 |
job_type | "batch"(默认)或"interactive" |
trigger | 单个或多个触发器(字符串或TTrigger),见下文触发器一节 |
execute | 执行约束:timeout(秒、"4h"之类人类可读串或TTimeoutSpec字典)、concurrency(最大并发运行数,默认1,传None取消限制) |
expose | UI 呈现:tags(分组标签)、starred(置顶)、manual(设False禁用手动触发) |
require | 运行时资源要求:dependency_groups、profile、instance(如{"size": "medium"})、region、static_egress_ips |
deliver | 关联的@dlt.source、独立@dlt.resource或已调用 source 实例 |
interval | 区间调度的整体时间范围 |
freshness | 上游新鲜度约束(单个字符串、TFreshnessConstraint或列表) |
incremental_mode | "interval"(增量区间由调度器管理)或"pipeline"(增量在 pipeline 内自持状态) |
refresh_propagation | "auto"(默认,透传上游刷新信号)/"always"(每次成功都清空下游prev_completed_run)/"block"(阻断传播) |
auto_refresh_pipeline_mode | 刷新运行时应用到作业内所有 pipeline 的刷新模式 |
spec | 可选的配置 spec 类 |
从实现看,被装饰函数被包装成JobFactory(decorators.py):它会用with_config注入配置(配置段为jobs.<section>.<name>)、保留原函数签名与类型(ParamSpec/TypeVar),并暴露.success/.fail/.completed/.is_fresh等属性;to_job_definition()把所有元数据序列化成语义化的TJobDefinition清单字典(入口点、触发器、执行规格、freshness、interval、require 等),这就是后续dlthub deploy生成部署清单的数据来源。
dlthub deploy与清单对账(reconciliation)
dlthub deploy是清单式部署的核心命令,执行流程(deployments.md):
- 导入
__deployment__.py收集所有作业; - 生成部署清单(描述每个作业的触发器、入口点、元数据的 JSON 文档);
- 把代码与配置同步到平台;
- 提交清单执行对账(reconciliation)。
对账结果:
| 状态 | 含义 |
|---|---|
| added | 新作业,将被创建 |
| updated | 作业定义发生变化,将被更新 |
| unchanged | 无变化,保持原样 |
| archived | 上次清单有、这次没有——触发器被禁用但运行历史保留 |
也就是说,从__deployment__.py中移除作业不会删除它,只会归档,保留运行历史与日志。部署前可以先预览:
# 只预览将发生的变化,不实际应用 dlthub deploy --dry-run # 导出完整展开后的清单为 YAML dlthub deploy --show-manifest第四步:定时调度——给装饰器传trigger=
让作业由 cron 驱动,只需给装饰器加一个trigger=:
from dlt.hub.run import trigger @run.pipeline( "ingest_breweries", trigger=trigger.schedule("0 * * * *"), # 每小时整点执行 ) def ingest_breweries(): ...改完再执行一次uv run dlthub deploy,触发器即被推送到平台,调度器随后按 cron 运行该作业。
基础触发器工厂
trigger.py 中的工厂函数完整列表(并与 triggers.md 的语义表对应):
| 写法 | 含义 |
|---|---|
trigger.every("5m") | 固定周期重复("5m"、"6h",也接受秒数 float) |
trigger.schedule("0 * * * *") | cron 表达式 |
trigger.once("2026-12-31T23:59:59Z") | 在指定时间戳执行一次(接受 ISO 字符串、datetime、date、unix 时间戳) |
"*/5 * * * *" | 裸 cron 字符串,自动检测 |
upstream_job.success | 上游作业成功完成后触发 |
upstream_job.fail | 上游作业失败后触发 |
upstream_job.completed | 上游成功或失败都触发(内部即(success, fail)元组) |
此外还有trigger.http(port, path)(交互式作业的 HTTP 触发器)、trigger.deployment()(代码部署后触发)、trigger.webhook(path)、trigger.tag(name)(标签广播)等工厂。
:::tip 从平台调整调度 也可以在 dltHub 平台的Manage Schedule对话框中直接修改 cron,适合临时暂停或微调而无需重新部署。若要让修改永久生效,请回到装饰器修改并运行dlthub deploy——装饰器始终是调度的唯一事实来源。 :::
高级特性一:多触发器与run_context
一个作业可以挂任意数量的触发器,在函数体内用注入的run_context["trigger"]区分是哪一个触发的:
from dlt.hub.run import TJobRunContext @run.job( trigger=[ trigger.schedule("0 * * * *"), upstream_ingest.success, ], ) def transform(run_context: TJobRunContext): if run_context["trigger"] == "schedule": ... elif run_context["trigger"] == "followup": ...TJobRunContext是 launcher 注入的字典,包含:run_id、trigger、refresh,以及调度器提供的interval_start/interval_end(见 triggers.md)。
高级特性二:follow-up 触发器——代码即依赖图
每个被装饰的作业都暴露.success、.fail、.completed三个触发器属性,用于把作业串成依赖图:
from dlt.hub.run import TJobRunContext @run.pipeline("transform_pipeline", trigger=ingest_job.success) def transform(run_context: TJobRunContext): ...follow-up 触发器的特点是上游一结束立即触发,无需轮询、无调度延迟(triggers.md)。从源码看,JobFactory.success/.fail/.completed直接映射到_triggers.job_success(job_ref)/job_fail(decorators.py),而job_success/job_fail工厂会把"name"、"section.name"、"jobs.section.name"三种格式的作业引用解析为规范 job_ref(trigger.py)。
高级特性三:调度器驱动的区间(Scheduler-driven intervals)
对于增量管道,用interval=声明整体时间范围,让平台为每次运行分配[interval_start, interval_end]窗口:
@run.pipeline( my_pipeline, interval={"start": "2026-01-01T00:00:00Z"}, trigger=trigger.schedule("*/3 * * * *"), ) def daily_ingest(run_context: TJobRunContext): start = run_context["interval_start"] end = run_context["interval_end"] # 把 start/end 传入 source,使其成为输入的纯函数 ...行为要点(triggers.md):
- 每次运行拿到刚流逝的区间;
- 错过的运行会自动补跑(backfill)——窗口持续向前延伸;
- 刷新时平台把区间指针重置回
interval.start; - 源码保持无状态——不需要游标持久化,也不需要查询 dlt state。
cron 与 every 的区间语义差异(triggers.md):
schedule触发器用 cron 表达式生成以绝对时刻为起止的区间。调度器在一个区间关闭时启动作业,交付刚流逝的窗口。例如每天 3 点执行的trigger.schedule("0 3 * * *"),5 月 26 日 3:00 启动的那次运行拿到的是 5 月 25 日 3:00 至 5 月 26 日 3:00 的区间。新部署的作业不会在部署时立即运行:5 月 25 日中午部署 3 点的作业,首次运行发生在 5 月 26 日 3:00。手动启动(如dlthub job trigger)时拿到的是"当前区间"——通常为空(interval_start == interval_end),增量作业应将其视为 no-op;唯一例外是错过的或失败的调度 tick,此时手动运行会补跑缺口。要强制重放已加载的窗口请用 refresh 而非手动运行。every触发器生成固定周期、从当前时刻起算的相对区间:14:20 部署trigger.every("1h"),首次运行在 15:20,区间为 14:20–15:20。手动运行时区间从上一次运行起点延伸到当前时刻,因此every 作业的手动运行总是拿到非空区间。
高级特性四:freshness 门控——上游未就绪就跳过
freshness=[upstream.is_fresh]会阻塞作业,直到上游最近一个区间完整结束:
@run.pipeline( "report_pipeline", trigger=trigger.schedule("0 * * * *"), freshness=[ingest_job.is_fresh], ) def build_report(run_context: TJobRunContext): ...与触发器不同:作业仍按自己的调度运行,只是在游上游加载进行中时跳过。适用于绝不能观察到部分数据的下游转换任务。注意当前 freshness 是单一水位线(只对上游最近完成的区间做门控),文档也提示区间级 freshness(每个到达的区间标记对应下游窗口过期)尚未支持(triggers.md)。源码中JobFactory.is_fresh与is_matching_interval_fresh对应_freshness.is_fresh/is_matching_interval_fresh两个约束构造器(decorators.py)。
高级特性五:刷新级联(Refresh cascade)
设置refresh_propagation="always"的补数作业会发起刷新信号,信号沿依赖图向下游所有作业传播;下游作业收到run_context["refresh"] = True后自行决定如何处理(例如pipeline.refresh = "drop_sources")。
刷新策略:
| 策略 | 行为 |
|---|---|
"always" | 每次运行都发起刷新信号 |
"auto" | 透传从上游收到的刷新信号(默认) |
"block" | 在此处阻断刷新传播 |
from dlt.hub import run @run.job(expose={"tags": ["backfill"]}, refresh_propagation="always") def backfill(): """级联刷新;不加载数据。"""用 CLI 触发:
dlthub job trigger "tag:backfill" dlthub run backfill --refresh # 对单个作业显式刷新刷新信号不会自动删除数据——需配合 dlt 的 pipeline.md 刷新选项 使用:
@run.pipeline( "report_pipeline", trigger=trigger.schedule("0 * * * *"), freshness=[ingest_job.is_fresh], ) def build_report(run_context: TJobRunContext): ... report_pipeline.run( data_source(), refresh="drop_data" if run_context["refresh"] else None )上面这段会在收到刷新信号时让 dlt 截断data_source()中资源所属的全部表(triggers.md)。
高级特性六:标签与批量触发
标签是挂在作业上的标签(通过expose={"tags": [...]}设置),用途有二:在仪表盘分组展示作业;在 CLI 用**选择器(selectors)**做批量操作:
# 触发所有打上 "ingest" 标签的作业 dlthub job trigger "tag:ingest" # 触发所有带调度计划的作业 dlthub job trigger "schedule:*" # 只预览不执行 dlthub job trigger "tag:ingest" --dry-run(示例见 triggers.md。)
高级特性七:时区
cron 表达式默认按UTC解释。要在特定 IANA 时区解释,在作业上声明:
@run.pipeline( my_pipeline, trigger=trigger.schedule("0 9 * * *"), # 早上 9 点 require={"timezone": "Europe/Berlin"}, # ……柏林时间 ) def morning_load(): ...run_context中的区间仍是 UTC datetime,但会与声明时区下的 tick 边界对齐。require={"timezone": ...}对整个运行生效,同时它也是该运行 dlt 的 context timezone:dlt 会以该时区读取 naive 时间戳、以该时区取date列的"天"。未声明时 dlt 使用 UTC(triggers.md)。
部署与配置独立版本化
在 dltHub 平台上,部署(你的代码文件)与配置(.dlt/*.toml)分开版本管理,可以只更新代码而不动 secrets,反之亦然:
# 只同步代码 / 只同步配置(不触发清单对账) dlthub workspace deployment sync dlthub workspace configuration sync # 查看历史版本 dlthub workspace deployment list dlthub workspace deployment info [version_number] dlthub workspace configuration list dlthub workspace configuration info [version]运行、监控与调试
部署后定时作业自动运行,也可手工触发;需要本地调试时先跑本地副本(deployments.md):
dlthub local run load_commits # 本地运行 dlthub run load_commits -f # 云端运行并跟随日志 dlthub job trigger "tag:ingest" # 不重新同步代码,用当前已部署代码触发监控命令(monitoring.md):
dlthub workspace info # workspace 概览:作业数、最新运行状态、版本 dlthub job list # 列出全部作业 dlthub job list "tag:ingest" # 按选择器过滤 dlthub job info <name> # 单个作业详情 dlthub job runs list [name_or_selector] [--running] dlthub job runs info <name> [run#] # 运行状态与耗时(默认最新一次) dlthub job logs my_pipeline.py 3 # 查看指定运行号的日志 dlthub job logs my_pipeline.py --follow # 实时流式日志 dlthub job runs cancel my_pipeline.py 5 # 取消指定运行 dlthub job cancel "tag:ingest" --dry-run # 按选择器批量取消(先预览)运行状态语义:Pending(排队等待)、Starting(初始化中)、Running(执行中)、Completed(无错完成)、Failed(出错,查日志)、Cancelled(手动停止)。Web UI 的 run 详情页提供状态栏、pipeline 运行表(每次作业执行过的 dlt pipeline 及其行数/状态)与实时日志查看器。
常见失败原因(monitoring.md):pyproject.toml缺依赖、prodprofile 的 secrets 未配置(平台批处理作业使用prod)、脚本缺if __name__ == "__main__":、残留dev_mode=True(每次运行重建 dataset)、prod与dev目标不一致、作业超时(默认 120 分钟,可用execute={"timeout": "6h"}覆盖)。
小结
dltHub 托管编排的完整工作流可以用一条链路概括:uvx dlthub-init@latest脚手架 →uv run dlthub run pipeline.py冒烟 → 加@run.pipeline装饰器并声明__deployment__.py→uv run dlthub deploy对账部署 →trigger=挂上 cron/间隔/follow-up → 用dlthub job系列命令监控与排障。调度的唯一事实来源始终是代码里的装饰器——这正是"代码即编排"的含义。需要进一步探索时,可继续阅读 triggers.md、deployments.md、job-configuration.md 与 monitoring.md,以及 command-line-interface.md 中deploy/run/job命令的完整参考。
【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy 🛠️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考