如何不使用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
相关产品推荐
相关产品推荐

