运行并跟踪工作流

使用 AI Studio SDK 启动已部署的工作流,并根据本次运行的状态读取结果或进行控制。启动后保存运行 ID;后续操作都以这个 ID 确定同一次运行。

任务流程

一次运行依次经过输入准备、提交启动、取得运行 ID、读取当前状态和读取结果几个阶段。运行 ID 是启动与后续查询、控制和结果读取之间的交接信息。

工作区 + 已部署工作流 + 运行输入
              │
              ▼
           启动运行
              │
              ▼
      运行 ID + 当前状态
        ├── 检查运行状态和可用操作
        ├── 读取运行结果
        └── 暂停、恢复、重试或取消

准备

需要的内容

如何取得

在本页中的作用

已认证的客户端

在首次接入流程中创建。

访问 AI Studio 资源。

工作区 ID

从工作区创建、列表或操作者选择中取得。

确定工作流的资源作用域。

已部署的工作流 ID

从部署工作流的结果或操作者选择中取得。

确定本次要启动的工作流。

工作流输入

从该工作流的输入定义和当前业务数据中取得。

为本次运行提供字段值。

本页不部署工作流,也不验证工作流输入定义。缺少这些内容时,先在相应的工作流管理流程中取得它们。

启动一次运行

提交运行后,SDK 返回本次运行的运行 ID 和当前状态。保存运行 ID;提交成功不表示工作流已经到达最终状态。

import os

import moi_product_sdk as sdk

client = sdk.new_with_personal_access_token(
    os.environ["PRODUCT_API_BASE_URL"],
    os.environ["PRODUCT_API_KEY"],
)
workspace = client.workspace("<workspace-id>")
workflow = workspace.workflow("<workflow-id>")
started = workflow.run(
    sdk.with_workflow_run_input(
        {"<input-field-id>": "<input-value>"}
    )
)

execution_id = started.workflow_run.execution_id
print(execution_id)
print(started.workflow_run.status)
package main

import (
	"context"
	"fmt"
	"os"

	sdk "github.com/matrixorigin/matrixflow/sdk/go-sdk"
)

func main() {
	ctx := context.Background()
	client, err := sdk.NewWithPersonalAccessToken(
		os.Getenv("PRODUCT_API_BASE_URL"),
		os.Getenv("PRODUCT_API_KEY"),
	)
	if err != nil {
		panic(err)
	}
	workspace, err := client.Workspace("<workspace-id>")
	if err != nil {
		panic(err)
	}
	workflow, err := workspace.Workflow("<workflow-id>")
	if err != nil {
		panic(err)
	}
	started, err := workflow.Run(
		ctx,
		sdk.WithWorkflowRunInput(map[string]any{
			"<input-field-id>": "<input-value>",
		}),
	)
	if err != nil {
		panic(err)
	}

	executionID := started.GetWorkflowRun().GetExecutionId()
	fmt.Println(executionID)
	fmt.Println(started.GetWorkflowRun().GetStatus())
}

检查本次运行

使用上一步的运行 ID 读取本次运行的当前状态、可用操作和错误信息。SDK 不提供通用的自动等待或自动重试循环;由应用根据返回状态决定后续频率和分支。

run = workflow.run_handle(execution_id)
current = run.refresh()

print(current.execution.status)
print(current.execution.available_actions)
print(current.execution.error)
run, err := workflow.RunHandle(executionID)
if err != nil {
	panic(err)
}
current, err := run.Refresh(ctx)
if err != nil {
	panic(err)
}

fmt.Println(current.GetExecution().GetStatus())
fmt.Println(current.GetExecution().GetAvailableActions())
fmt.Println(current.GetExecution().GetError())

读取运行结果

当应用需要结果记录时,继续使用同一运行引用读取结果。结果中包含本次运行的状态、业务结果、错误信息、节点状态和可用操作。整体运行的状态不能替代某个节点的状态。

result = run.result()

print(result.result.status)
print(result.result.case_result)
print(result.result.case_error)
result, err := run.Result(ctx)
if err != nil {
	panic(err)
}

fmt.Println(result.GetResult().GetStatus())
fmt.Println(result.GetResult().GetCaseResult())
fmt.Println(result.GetResult().GetCaseError())

控制运行

在状态和可用操作表明可继续控制时,可以暂停、恢复、重试或取消本次运行。重试会返回新的运行信息,因此应保存新运行 ID;其他控制操作返回当前运行信息,仍需读取返回状态或再次检查确认。

目标

操作结果

暂停或恢复

返回当前运行信息。

重试

返回新的运行信息。

取消

返回当前运行信息。

下一步

最后更新于