跳转到内容

可复用工作流

Workflow 保存一次混合计算的任务依赖。它把算法结构与运行时对象分开:定义时用别名引用组件和量子后端,运行时再提供当前 Runtime 中的对象。

依次调用 workflow.input() 声明输入,用 task()、component()、quantum() 添加节点,用 output() 声明需要取回的输出。每个节点只能引用已声明输入或较早的节点,因此图不会出现循环。

workflow = pq.Workflow("feedback")
theta = workflow.input("theta")
circuit = workflow.task(build_circuit, theta, name="prepare")
measured = workflow.quantum("quantum", circuit, shots=1024, name="measure")
updated = workflow.task(update_parameter, theta, measured, name="update")
workflow.output("updated", updated)
with pq.Runtime() as runtime:
quantum = runtime.quantum_backend("simulator")
run = runtime.run(workflow, inputs={"theta": 1.0}, bindings={"quantum": quantum})
result = runtime.get(run.outputs["updated"])
run.release()

上面的 build_circuit、update_parameter 是用户函数;下面的完整教程提供可直接运行的实现。NodeRef 表示图中的值,执行后 WorkflowRun.outputs 才包含可供 get() 使用的 ResultRef。

程序重复运行同一个工作流,用 Actor 累积量子结果,最后导出执行记录:

终端窗口
python packages/framework/examples/system_workflow.py --report-path execution.json
python packages/framework/examples/system_workflow.py --executor ray --address local

完整教程代码如下:

system_workflow.py
"""A reusable CPU/quantum workflow with a stateful actor and execution report.
Run: python examples/system_workflow.py
Ray: python examples/system_workflow.py --executor ray --address local
This example uses only a quantum simulator and is independent of AIMD.
"""
from __future__ import annotations
import argparse
import json
import pivotq as pq
def build_circuit(theta: float):
from pivotq import QuantumCircuit
circuit = QuantumCircuit(3)
circuit.ry(theta, 0)
circuit.cx(0, 1)
circuit.cx(1, 2)
circuit.measure_all()
return circuit
class History:
def __init__(self):
self.values = []
def record(self, measurement: pq.QuantumResult):
self.values.append(measurement.probabilities.get("111", 0.0))
return {"steps": len(self.values), "probabilities": list(self.values)}
def make_workflow():
workflow = pq.Workflow("quantum-feedback")
theta = workflow.input("theta")
circuit = workflow.task(build_circuit, theta, name="prepare")
measurement = workflow.quantum("quantum", circuit, shots=512, seed=7, name="measure")
recorded = workflow.component("history", "record", measurement, name="record")
workflow.output("history", recorded)
return workflow
def run(*, executor="local", address=None, report_path=None):
workflow = make_workflow()
with pq.Runtime(executor=executor, address=address, trace=True) as runtime:
history = runtime.actor(History, methods=("record",), name="history")
quantum = runtime.quantum_backend("simulator")
for theta in (0.5, 1.0):
execution = runtime.run(workflow, inputs={"theta": theta},
bindings={"history": history, "quantum": quantum})
ready, pending = runtime.wait(execution.refs, num_returns=len(execution.refs), timeout=60)
if pending:
raise RuntimeError("example workflow did not finish before the wait timeout")
result = runtime.get(execution.outputs["history"])
assert runtime.status(execution.outputs["history"]) is pq.InvocationStatus.SUCCEEDED
execution.release()
report = runtime.report()
if report_path is not None:
report.export(report_path)
return {"result": result, "trace_records": len(report.records),
"report_available_after_close": report.closed,
"is_simulated": True, "dropped_records": report.dropped_records}
def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--executor", choices=("local", "ray"), default="local")
parser.add_argument("--address")
parser.add_argument("--report-path")
print(json.dumps(run(**vars(parser.parse_args())), indent=2))
if __name__ == "__main__":
main()

第一次提交时工作流结构会冻结。之后可更换输入和绑定再次运行,每次得到独立的引用。绑定的 ComponentHandle、QuantumBackend 必须属于运行该图的 Runtime;输入和绑定名称必须与声明匹配。

运行时会预检结构和接口。执行系统在实际提交过程中仍可能发生故障;如果部分节点已被接收,WorkflowSubmissionError.partial_run 会保留这些节点的引用。此时不能假定整个工作流都未执行,并据此重新提交。

run.refs 包含该次执行所有节点的引用。先用 runtime.wait(run.refs, num_returns=len(run.refs)) 等待需要释放的节点完成,再调用 run.release()。只有一个分支的输出已完成,不表示其他分支也已完成。

算法中的 for、while 循环,以及依赖结果的条件分支,仍用普通 Python 编写。工作流描述执行依赖;性能工作量模型是另一个由用户编写的模型,不会自动从图中推断准确计算耗时。