workflow 教程
本教程展示如何把并发工作描述为任务图:声明能力,添加使用它们的节点,用边为节点排序,并在任何东西运行之前让 validate 检查结果。这个包可以在所有目标上使用。
快速开始
导入该包;需要策略和计算节点时,再导入 shared 和 plan:
import {
"Luna-Flow/luna_thread/plan",
"Luna-Flow/luna_thread/shared",
"Luna-Flow/luna_thread/workflow",
}
一个只有两个节点的派生-汇合图:
test "workflow quick start" {
let graph = @workflow.Workflow::new("fork-join")
.add_node(@workflow.spawn_node(1, "spawn"))
.add_node(@workflow.join_node(2, "join"))
.add_edge(@workflow.Edge::new(1, 2, @workflow.control_dependency()))
inspect(@workflow.is_ready(graph), content="true")
}
日常任务
通过通道传递工作
相互通信的节点通过 id 指定能力。通道以 MoveOnly 访问模式使用,从 send 到 recv 的边使接收等待发送完成:
test "channel" {
let jobs = @workflow.Capability::new(
1,
"jobs",
@workflow.channel_capability(),
@workflow.move_only_access(),
)
let graph = @workflow.Workflow::new("producer-consumer")
.add_capability(jobs)
.add_node(@workflow.send_node(1, "produce", capability=1))
.add_node(@workflow.recv_node(2, "consume", capability=1))
.add_edge(@workflow.Edge::new(1, 2, @workflow.data_dependency()))
assert_eq(@workflow.validate(graph).length(), 0)
}
保护临界区
互斥锁以 SynchronizeOnly 访问模式,由一个 Lock 节点和一个 Unlock 节点使用:
test "critical section" {
let guard_cap = @workflow.Capability::new(
1,
"guard",
@workflow.mutex_capability(),
@workflow.synchronize_only_access(),
)
let graph = @workflow.Workflow::new("critical")
.add_capability(guard_cap)
.add_node(@workflow.lock_node(1, "enter", capability=1))
.add_node(@workflow.write_shared_node(2, "update", capability=2))
.add_node(@workflow.unlock_node(3, "leave", capability=1))
.add_edge(@workflow.Edge::new(1, 2, @workflow.synchronization_dependency()))
.add_edge(@workflow.Edge::new(2, 3, @workflow.synchronization_dependency()))
debug_inspect(@workflow.validate(graph), content="[MissingCapability(capability_id=2)]")
}
写入节点指定了能力 2,而图中没有声明它。把一个访问模式为 ReadWrite 的 AtomicCell 作为能力 2 加入即可修正。
嵌入计划
计算节点携带一个计划,计划自身的问题会连同节点 id 一起报告:
test "compute node" {
let plan = @plan.map("halve", @plan.f64_type(), 8)
let graph = @workflow.Workflow::new("compute").add_node(
@workflow.compute_node(1, "halve", plan),
)
debug_inspect(
@workflow.validate(graph),
content="[InvalidComputePlan(node_id=1, issue=UnsupportedValueType)]",
)
assert_true(@workflow.validate(graph).any(@workflow.is_invalid_compute_plan))
}
提交工作流
submit 验证图,并记录所选后端是否接受它:
test "submit" {
let graph = @workflow.Workflow::new("one").add_node(
@workflow.spawn_node(1, "spawn"),
)
let submission = @workflow.submit(graph)
inspect(submission.accepted(), content="true")
inspect(submission.completed(), content="false")
}
进阶
要运行图,请在原生目标上把它传给 @native.submit_workflow_async(或门面的 submit_workflow_async);原生后端教程展示了做法,原生后端设计解释了哪些节点种类会阻塞以及什么会唤醒它们。由于 validate 返回所有问题,工具可以一次展示图的所有问题;谓词 is_missing_capability、is_cyclic_dependency、is_invalid_compute_plan 和 is_unsupported_shared_write 对常见问题做了分组。
常见陷阱
add_capability、add_node和add_edge会就地修改工作流。绑定到同一个工作流的两个变量看到的是同样的节点。- 当某个节点没有入边时,
validate会漏掉三个或更多节点构成的环;原生运行时会在提交时拒绝它们。 - 可能先于其配对节点运行的
Recv、Lock或Send可能会在原生运行时上永远阻塞。请用边为这类节点排序。 - 组内只有一个
Barrier节点时,会在运行时以状态码14失败;一个屏障需要在同一深度上至少有两个节点。 RwLock和Semaphore能通过validate,但会被原生运行时拒绝。
下一步
workflow API 列出了 validate 的每项检查。workflow 设计给出了图模型和环检查。core 教程通过门面构建工作流。