ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

使用 dltHub 托管编排部署 dlt 数据管道:从脚手架、部署到定时触发与刷新级联

使用 dltHub 托管编排部署 dlt 数据管道:从脚手架、部署到定时触发与刷新级联 使用 dltHub 托管编排部署 dlt 数据管道从脚手架、部署到定时触发与刷新级联【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dltdltHub 是 dltdata 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 并脚手架化一个 workspacedltHub 的 CLI 名为dlthub通过uvx直接拉起最新版本即可无需手工安装uvx dlthub-initlatest如果你还没有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_nameingest_breweries, destinationwarehouse, dataset_namebrewery_data, ) pipeline.run(brewery_source())__deployment__.pyworkspace 的部署清单# __deployment__.py from ingest_breweries import ingest_breweries __all__ [ingest_breweries]然后部署uv run dlthub deploy之后按作业名触发运行uv run dlthub run ingest_breweries装饰器三兄弟pipeline / job / interactivedlt/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_typebatch默认或interactivetrigger单个或多个触发器字符串或TTrigger见下文触发器一节execute执行约束timeout秒、4h之类人类可读串或TTimeoutSpec字典、concurrency最大并发运行数默认1传None取消限制exposeUI 呈现tags分组标签、starred置顶、manual设False禁用手动触发require运行时资源要求dependency_groups、profile、instance如{size: medium}、region、static_egress_ipsdeliver关联的dlt.source、独立dlt.resource或已调用 source 实例interval区间调度的整体时间范围freshness上游新鲜度约束单个字符串、TFreshnessConstraint或列表incremental_modeinterval增量区间由调度器管理或pipeline增量在 pipeline 内自持状态refresh_propagationauto默认透传上游刷新信号/always每次成功都清空下游prev_completed_run/block阻断传播auto_refresh_pipeline_mode刷新运行时应用到作业内所有 pipeline 的刷新模式spec可选的配置 spec 类从实现看被装饰函数被包装成JobFactorydecorators.py它会用with_config注入配置配置段为jobs.section.name、保留原函数签名与类型ParamSpec/TypeVar并暴露.success/.fail/.completed/.is_fresh等属性to_job_definition()把所有元数据序列化成语义化的TJobDefinition清单字典入口点、触发器、执行规格、freshness、interval、require 等这就是后续dlthub deploy生成部署清单的数据来源。dlthub deploy与清单对账reconciliationdlthub deploy是清单式部署的核心命令执行流程deployments.md导入__deployment__.py收集所有作业生成部署清单描述每个作业的触发器、入口点、元数据的 JSON 文档把代码与配置同步到平台提交清单执行对账reconciliation。对账结果状态含义added新作业将被创建updated作业定义发生变化将被更新unchanged无变化保持原样archived上次清单有、这次没有——触发器被禁用但运行历史保留也就是说从__deployment__.py中移除作业不会删除它只会归档保留运行历史与日志。部署前可以先预览# 只预览将发生的变化不实际应用 dlthub deploy --dry-run # 导出完整展开后的清单为 YAML dlthub deploy --show-manifest第四步定时调度——给装饰器传trigger让作业由 cron 驱动只需给装饰器加一个triggerfrom dlt.hub.run import trigger run.pipeline( ingest_breweries, triggertrigger.schedule(0 * * * *), # 每小时整点执行 ) def ingest_breweries(): ...改完再执行一次uv run dlthub deploy触发器即被推送到平台调度器随后按 cron 运行该作业。基础触发器工厂trigger.py 中的工厂函数完整列表并与 triggers.md 的语义表对应写法含义trigger.every(5m)固定周期重复5m、6h也接受秒数 floattrigger.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, triggeringest_job.success) def transform(run_context: TJobRunContext): ...follow-up 触发器的特点是上游一结束立即触发无需轮询、无调度延迟triggers.md。从源码看JobFactory.success/.fail/.completed直接映射到_triggers.job_success(job_ref)/job_faildecorators.py而job_success/job_fail工厂会把name、section.name、jobs.section.name三种格式的作业引用解析为规范 job_reftrigger.py。高级特性三调度器驱动的区间Scheduler-driven intervals对于增量管道用interval声明整体时间范围让平台为每次运行分配[interval_start, interval_end]窗口run.pipeline( my_pipeline, interval{start: 2026-01-01T00:00:00Z}, triggertrigger.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.mdschedule触发器用 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, triggertrigger.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_propagationalways的补数作业会发起刷新信号信号沿依赖图向下游所有作业传播下游作业收到run_context[refresh] True后自行决定如何处理例如pipeline.refresh drop_sources。刷新策略策略行为always每次运行都发起刷新信号auto透传从上游收到的刷新信号默认block在此处阻断刷新传播from dlt.hub import run run.job(expose{tags: [backfill]}, refresh_propagationalways) def backfill(): 级联刷新不加载数据。用 CLI 触发dlthub job trigger tag:backfill dlthub run backfill --refresh # 对单个作业显式刷新刷新信号不会自动删除数据——需配合 dlt 的 pipeline.md 刷新选项 使用run.pipeline( report_pipeline, triggertrigger.schedule(0 * * * *), freshness[ingest_job.is_fresh], ) def build_report(run_context: TJobRunContext): ... report_pipeline.run( data_source(), refreshdrop_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, triggertrigger.schedule(0 9 * * *), # 早上 9 点 require{timezone: Europe/Berlin}, # ……柏林时间 ) def morning_load(): ...run_context中的区间仍是 UTC datetime但会与声明时区下的 tick 边界对齐。require{timezone: ...}对整个运行生效同时它也是该运行 dlt 的 context timezonedlt 会以该时区读取 naive 时间戳、以该时区取date列的天。未声明时 dlt 使用 UTCtriggers.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.mddlthub local run load_commits # 本地运行 dlthub run load_commits -f # 云端运行并跟随日志 dlthub job trigger tag:ingest # 不重新同步代码用当前已部署代码触发监控命令monitoring.mddlthub 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.mdpyproject.toml缺依赖、prodprofile 的 secrets 未配置平台批处理作业使用prod、脚本缺if __name__ __main__:、残留dev_modeTrue每次运行重建 dataset、prod与dev目标不一致、作业超时默认 120 分钟可用execute{timeout: 6h}覆盖。小结dltHub 托管编排的完整工作流可以用一条链路概括uvx dlthub-initlatest脚手架 →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),仅供参考
返回列表