You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何不使用Dagit GraphQL以编程方式提交Dagster运行?关于pipeline_run与workspace获取及替代方案咨询

没问题,我来帮你梳理下不用Dagit GraphQL API提交Dagster运行的几种靠谱方式,顺便解决你当前代码里的问题:

首先纠正你当前代码的问题

你之前的代码里误用了submit_run的参数:workspace=graphene_info.context是Dagit内部用来处理Web请求的上下文,完全不需要在编程提交时传入;而且pipeline_run需要你手动构建,不是直接获取的。


方法1:用DagsterInstance.submit_run异步提交运行

这是最接近你最初思路的方式,直接和Dagster实例交互,异步提交运行到后台队列(如果配置了Celery等执行器)。

完整代码示例

from dagster import DagsterInstance, PipelineRun, RunConfig
import uuid

# 获取本地Dagster实例(默认读取~/.dagster/dagster.yaml配置)
instance = DagsterInstance.get()

# 生成唯一运行ID(也可以不传,让实例自动生成)
run_id = str(uuid.uuid4())

# 构建PipelineRun对象,核心是指定你的流水线名称和配置
pipeline_run = PipelineRun(
    pipeline_name="your_pipeline_name",  # 替换成你的流水线名称
    run_id=run_id,
    run_config=RunConfig(
        # 这里传入你的流水线配置,比如给solid传参数
        solids={"your_solid_name": {"config": {"input_value": "test"}}}
    ),
    mode="default"  # 流水线的运行模式,默认是"default"
)

# 提交运行
submitted_run = instance.submit_run(pipeline_run)

print(f"成功提交运行,ID: {submitted_run.run_id}")

关键说明

  • 不需要workspace参数:DagsterInstance会自动读取你配置的仓库(比如dagster.yaml里的code_locations),只要你的流水线已经注册到实例中就能找到。
  • 如果是多仓库场景,你可以通过instance.get_repository_location("your_location_name")先定位仓库,再获取流水线,但一般单仓库场景直接用PipelineRun指定名称即可。

方法2:用execute_pipeline同步执行流水线

如果不需要异步提交,而是想直接在脚本里同步运行并获取结果,这个方法更简单,完全脱离Dagit依赖。

完整代码示例

from dagster import execute_pipeline, RunConfig
# 导入你定义的流水线
from your_pipeline_module import your_pipeline_name

# 同步执行流水线,会阻塞到运行完成
run_result = execute_pipeline(
    your_pipeline_name,
    run_config=RunConfig(
        solids={"your_solid_name": {"config": {"input_value": "test"}}}
    ),
    mode="default"
)

# 检查运行结果
if run_result.success:
    print("流水线运行成功!")
    # 可以获取输出数据
    output_data = run_result.output_for_solid("your_solid_name")
    print(f"Solid输出: {output_data}")
else:
    print("流水线运行失败,查看日志获取详情")

关键说明

  • 适合测试、脚本化运行场景,不需要启动Dagit服务,直接运行脚本即可。
  • 可以直接获取运行结果和Solid的输出,方便后续处理。

方法3:调用Dagster CLI(适合集成/自动化场景)

如果你的场景更倾向于用命令行工具集成,也可以在Python里通过subprocess调用Dagster的CLI命令来提交运行。

完整代码示例

import subprocess

# 调用dagster CLI提交运行
subprocess.run(
    [
        "dagster", "pipeline", "execute",
        "--pipeline-name", "your_pipeline_name",
        # 指定配置文件路径,或者用--config-yaml直接传配置字符串
        "--config", "path/to/your_run_config.yaml",
        # 如果你的仓库在自定义位置,指定仓库yaml文件
        "--repository-yaml", "path/to/repository.yaml"
    ],
    check=True,  # 如果命令执行失败会抛出异常
    capture_output=True,
    text=True
)

print("流水线提交成功")

关键说明

  • 适合已经熟悉Dagster CLI的场景,不需要写太多Dagster Python API代码。
  • 可以轻松集成到shell脚本、CI/CD流程中。

内容的提问来源于stack exchange,提问作者Varun Shridhar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.28 20:17:43