使用 dltHub 托管编排部署 dlt 数据管道:从脚手架、部署到定时触发与刷新级联
2026/9/18 14:16:25 网站建设 项目流程

使用 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导出jobpipelineinteractive三个装饰器、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取消限制)
exposeUI 呈现:tags(分组标签)、starred(置顶)、manual(设False禁用手动触发)
require运行时资源要求:dependency_groupsprofileinstance(如{"size": "medium"})、regionstatic_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):

  1. 导入__deployment__.py收集所有作业;
  2. 生成部署清单(描述每个作业的触发器、入口点、元数据的 JSON 文档);
  3. 把代码与配置同步到平台;
  4. 提交清单执行对账(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_idtriggerrefresh,以及调度器提供的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_freshis_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)、proddev目标不一致、作业超时(默认 120 分钟,可用execute={"timeout": "6h"}覆盖)。

小结

dltHub 托管编排的完整工作流可以用一条链路概括:uvx dlthub-init@latest脚手架 →uv run dlthub run pipeline.py冒烟 → 加@run.pipeline装饰器并声明__deployment__.pyuv 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询