Workflow composition, failure handlers, and nodes
从函数到 DAG
当你在普通 Python 函数上使用 @workflow 进行装饰时,flytekit 会将其包装在 PythonFunctionWorkflow 中(flytekit/core/workflow.py)。与其他按需运行的代码不同,工作流函数体在序列化时会执行一次:flytekit 会运行该函数,但会拦截函数内部的每一次任务调用,并将其转换为 DAG 节点,而不是运行代码。
输入是如何变为 promise 的:在编译期间,flytekit 会使用 workflow.py(第 844 行)中 compile() 方法内的 construct_input_promises 构建一个 promise 字典。每个工作流输入都会变成一个绑定到 GLOBAL_START_NODE 的 Promise —— 这是一个在模块顶部定义的单例输入节点,表示整个工作流的输入。然后使用这些 promise 调用该函数,因此当函数体运行时,每个变量都是指代工作流输入节点的输出的 Promise。
任务调用是如何变为节点的:当函数体调用像 t1(a=a) 这样的任务时,调用会被 flytekit/core/promise.py 中的 flyte_entity_call_handler 拦截。在编译过程中,这会创建一个新的 Node(来自 flytekit/core/node.py),将其附加到当前的 CompilationState,并返回包装了 NodeOutput 引用的 Promise 对象,而不是实际结果。返回的 promise 会变成工作流的输出绑定。
因为工作流函数体仅在序列化时运行,所以内部传递的值并不是普通的 Python 值 —— 它们是 Promise 对象。如果 t1() -> int,那么 a = t1() 不会产生整数;如果你尝试对 a 进行 range() 或对其进行真值测试,flytekit 会引发异常(例如,Promise.__bool__ 会引发 "Flytekit does not support Unary expressions or performing truth value testing...")。
比较运算符(==、<、> 等)会被重载以产生用于 conditional 表达式的 ComparisonExpression 对象,而不是布尔值。逻辑 and / or 不能被重载,因此必须使用位运算符 & 和 |,它们会产生 ConjunctionExpression 对象。在 promise 上调用 int() 或 float() 等内置转换函数也会失败,因为会调用 __bool__。
from flytekit import task, workflow, conditional
@task
def check(x: int) -> bool:
return x > 5
@task
def hi() -> str:
return "high"
@task
def lo() -> str:
return "low"
@workflow
def wf(x: int) -> str:
return conditional("branch").if_(check(x=x)).then(hi()).else_().then(lo())
create_node:直接访问节点级输出
调用 create_node(entity, ...)(来自 flytekit/core/node_creation.py)是使用节点而不是仅使用 promise 的方法。它会调用实体(类似于调用任务),然后从编译状态中获取新创建的 Node 并将输出 promise 存储在其上:
node._outputs = {}
# ... for each output:
setattr(node, output_name, attr)
node.outputs[output_name] = attr
这意味着返回的节点会公开输出两种方式:作为属性(node.o0,匹配任务声明的输出名称)和通过 node.outputs 字典(键为输出名称)。
一个关键的区别:create_node(...).outputs 是节点级输出访问器 —— 它存在于 Node 对象本身,允许你重用同一个节点的输出,而无需多次调用该任务。相比之下,像 t1(a=x) 这样直接调用任务会立即返回 Promise 对象(或命名元组),而不会保留对节点的直接句柄。
例如,从任务的 docstring 来看,对于一个 t4() -> (int, str):
t4_node = create_node(t4)
# In compilation node.o0 has the promise.
t5(in1=t4_node.o0)
如果你需要在下游通过名称引用输出,这也很有用:
t3_node = create_node(t3, in1=some_int)
这基本上就是 t3(in1=some_int),但你会得到 Node,以便以后可以访问 t3_node.outputs["out1"] 或将其用于显式连线。
请注意 docstring 中关于本地执行的注意事项:即使在本地执行中,对于单输出任务,你仍然会得到一个需要通过输出名称解引用的包装器:
t1_node = create_node(t1)
t2(t1_node.o0)
指定不产生输出的任务之间的依赖关系
create_node 对于在既不消耗也不产生输出的任务之间表达依赖关系至关重要 —— 否则无法进行连线。从 docstring 中可以看出:
t1_node = create_node(t1)
t2_node = create_node(t2)
t2_node.runs_before(t1_node)
# OR
t2_node >> t1_node
Node 上的 >> 操作符会调用 runs_before,它会将 self 添加到 other 的 upstream_nodes 列表中。没有对应的重载来表示相反方向的依赖关系。
工作流中的条件执行
conditional 构造利用了重载的比较运算符。当你在 if_() 子句中比较 promise(例如 check(x=x) 的结果)时,flytekit 会构建一个包含 ComparisonOps(EQ、NE、GT、GE、LT、LE)的 ComparisonExpression。这些可以通过 & 和 | 组合成 ConjunctionExpression(带有 ConjunctionOps.AND / OR),其求值方式是惰性的/在编译时确定的。
失败处理程序
@workflow 上的 on_failure 参数指定了在工作流失败时调用的任务或工作流。失败实体的接口必须与工作流接口的超集兼容:
- 它必须接受工作流的每个输入。
- 它可以声明 额外的 输入,但每个额外输入必须是
Optional(通常是一个用于接收失败详细信息的err)。
验证在编译时进行。对于函数式工作流,PythonFunctionWorkflow 中的 _validate_add_on_failure_handler 会检查工作流输入是否为失败节点输入的子集,并对每个额外键调用 is_optional_type,如果违反约束则会引发 FlyteFailureNodeInputMismatchException。命令式工作流在调用 add_on_failure_handler 时也会执行相同的检查。
满足约束的有效示例:
@task
def clean_up(name: str, err: typing.Optional[FlyteError] = None):
print(f"Deleting cluster {name} due to {err}")
@workflow(on_failure=clean_up)
def wf(name: str = "flyteorg"):
c = create_cluster(name=name)
t = t1(a=1, b="2")
d = delete_cluster(name=name)
c >> t >> d
这里 clean_up 接受 name(存在于工作流上)和可选的 err(附加输入)。在执行期间,如果节点失败,工作流的 __call__ 会捕获异常,并且如果失败处理程序声明了 err 输入,flytekit 会构建一个 FlyteError(failed_node_id=..., message=str(exc)) 并将其传递给失败处理程序。
不符合约束的示例 —— bad_handler 声明了非可选的 err: FlyteError 而没有默认值,因此 Optional 检查失败:
@task
def bad_handler(name: str, err: FlyteError): # NOT optional -> FlyteFailureNodeInputMismatchException
...
同样,遗漏工作流输入(例如,仅接受 err 而不接受 name)会违反超集检查。