ワークフローと実行例
ワークフローの操作について
ここでは、ワークフローを以下の方法で操作します。
PREFECT_API_URLの設定
ワークフロープログラムの動作には環境変数 PREFECT_API_URL の設定が必要です。
Prefectサーバを起動した計算ノードのIPアドレスを指定します。
Prefectサーバとは異なる計算ノードでワークフローを操作する場合は、忘れずに設定してください。
ここでは、IPアドレス10.0.0.33の計算ノードでPrefectサーバが起動されたとします。
(venv_name) [username@qh001 ~]$ export PREFECT_API_URL='http://10.0.0.33:4200/api'
実行
ターミナルで、以下のコマンドを実行します。
(venv_name) [username@qh001 ~]$ python <workflow_program_file>
workflow_program_file: ワークフロープログラムファイルのパス
停止
Prefect のワークフローを停止するためには deployment という仕組みを利用する必要があります。
停止操作は、deployment を行ったワークフローでのみ可能です。
Info
deployment については、Prefect のドキュメントのDeploymentsを参照ください。
ここではコマンドラインで操作していますが、Prefect UI からでも停止することは可能です。
Note
ワークフローの停止は、実行中のターミナルから Ctrl + C を入力することでも可能です。
その際、タスクの状態は Cancelled ではなく、Crashed となります。また、Prefect UI でワークフロー図を正しく表示できない場合があります。
ワークフローの停止は、以下の流れで行います。
ワークフロープログラムの作成
タスクを停止できることを確認するためのワークフロープログラムを作成します。 タスクを停止できることを確認するためのワークフロープログラムの作成方法については、Prefect のドキュメントのHow to create deploymentsを参照ください。
deploymentの実行
作成したワークフローを deployment として登録します。
ターミナル上から以下のコマンドを実行します。
(venv_name) [username@qh001 ~]$ python <workflow_program_file> &
workflow_program_file: タスクを停止できることを確認するためのワークフロープログラム
上記コマンドを実行すると、以下のようなメッセージがターミナルに表示されます。
Your flow '<flow_name>' is being served and polling for scheduled runs!
To trigger a run for this flow, use the following command:
$ prefect deployment run '<flow_name>/<deployment_name>'
You can also run your flow via the Prefect UI: http://<ip_address>:<port>/deployments/deployment/a2476573-575c-4367-b7ff-58d61206c543
flow_name: ワークフロープログラムに記述したフローメソッド名
deployment_name: ワークフロープログラムのflow_name.serve メソッドで指定した deployment 名
ip_address: Prefectサーバーを起動する計算ノードのIPアドレス
port: ポート番号
ワークフローの実行
ワークフローを実行するために、ターミナル上で提示されたコマンドを実行します。
(venv_name) [username@qh001 ~]$ prefect deployment run '<flow_name>/<deployment_name>'
ワークフローの停止
実行したターミナル上から以下のコマンドを実行し、ワークフローを停止させます。
(venv_name) [username@qh001 ~]$ prefect flow-run cancel <UUID>
UUID: ワークフロー実行時に表示される UUID
監視
Prefect UI にアクセスして、ワークフロー図を表示させることで、タスクの実行状態を監視できます。
システム H でタスク実行
システム H でのタスクの実行方法について説明します。
ワークフロープログラムの作成
システム H で PBS を介してタスク実行するワークフロープログラムを作成します。
ここでは、フィボナッチ数列を3回計算して単純に足し合わせるワークフロープログラムの例を示します。 Fib(5) = 5 なので、答えは 15 になる想定です。
dask-jobqueue には dask_jobqueue.PBSCluster というPBS向けのクラスがあるので、それを使用して PrefectのワークフローにPBS用の設定を追記します。
PBS用に追記する内容は、以下の4つです。
cluster_kwargsにdask_jobqueue.PBSClusterの設定を記入dask_jobqueue.PBSClusterのインスタンス作成時の引数としてcluster_kwargsを指定dask_jobqueue.PBSClusterのscaleメソッドの引数で、PBSのジョブ数を指定@flowデコレーターの引数に、task_runner=DaskTaskRunnerを指定
ここでは、各ノードあたり 1 つの Worker を起動し、 2 ノードでタスク実行を行うワークフロープログラムとなっています。
dask_jobqueue.PBSCluster の主なパラメータ |
説明 |
|---|---|
| cores | CPU のコア数 |
| job_script_prologe | worker を起動するPBSジョブスクリプトのプロローグ処理。module load コマンド実行が必要です |
| walltime | ジョブの時間制限 |
| n_workers | worker の数 |
| local_directory | dask の worker のローカルディレクトリ |
| resource_spec | ABCI-Q の計算資源。ここでは rt_QF を指定し、占有ノードで実行させています |
| job_extra_directives | PBSジョブの起動時の引数。ABCI-Q のグループ名指定が必要です |
Info
dask_jobqueue.PBSCluster に設定できるパラメータについては、dask_jobqueue.PBSCluster のドキュメントを参照ください。
次のプログラムをtest_sysH_task.pyとして保存します。
プログラム内の<group_name>は所属するグループ名に置き換えてください。
from dask.distributed import Client, LocalCluster
from dask_jobqueue import PBSCluster
from prefect import flow, task
from prefect_dask import DaskTaskRunner
cluster_kwargs = {
"cores": 1,
"memory": "1 GB",
"job_script_prologue" : ["source /etc/profile.d/modules.sh","module load python/3.11/3.11.14"],
"walltime": "00:05:00",
"n_workers" : 1,
"local_directory" : "~/tmp",
"resource_spec" : "rt_QF=1",
"job_extra_directives" : ["-W group_list=<group_name>"]
}
cluster = PBSCluster(**cluster_kwargs)
pbs_jobs = 2
cluster.scale(jobs=pbs_jobs)
client = Client(cluster)
def fib(n: int) -> int:
if n <= 2:
return 1
return fib(n-1) + fib(n-2)
@task
def task1(_: int) -> int:
return fib(5)
@task
def task_sum(ls: list[int]) -> int:
ret = 0
for x in ls:
ret += x
return ret
@flow(task_runner=DaskTaskRunner(address=client.scheduler.address), log_prints=True)
def main():
xs = []
for i in range(3):
xs.append(task1.submit(i))
print(task_sum.submit(xs).result())
if __name__ == "__main__":
main()
Info
ワークフロープログラムを実行すると、PBS を介した dask worker の起動処理が行われ、その後 dask によってタスクが実行されます。タスクが完了すると、dask worker 起動のための PBS ジョブが kill されて、dask worker も kill されるという処理のフローとなります。
ワークフローの実行
作成したワークフロープログラムを実行すると、以下のようになります。
想定される 15という結果が得られていることがわかります。
また、実行したワークフローは Prefect UI から確認することができます。
(venv_name) [username@qh001 ~]$ python test_sysH_task.py
15:05:28.486 | INFO | Flow run 'cocky-rat' - Beginning flow run 'cocky-rat' for flow 'main'
15:05:28.489 | INFO | Flow run 'cocky-rat' - View at http://10.0.0.33:4200/runs/flow-run/cb84172b-9ce2-45fb-9159-6d05223e43e4
15:05:28.490 | INFO | prefect.task_runner.dask - Connecting to existing Dask cluster PBSCluster(e2325e4e, 'tcp://10.0.18.52:41785', workers=0, threads=0, memory=0 B)
15:05:50.701 | INFO | Flow run 'cocky-rat' - 15
15:05:50.716 | INFO | Flow run 'cocky-rat' - Finished in state Completed()
(venv_name) [username@qh001 ~]$
ワークフローの停止
Prefect のワークフローを停止するためには deployment という仕組みを利用する必要があります。
ワークフロープログラムの作成
システム H で実行したタスクを停止できることを確認するためのワークフロープログラムを作成します。フロー main() に対して serve というメソッドを呼び出すことで deployment を作成するようにしています。
次のプログラムをtest_sysH_deployment.pyとして保存します。
プログラム内の<group_name>は所属するグループ名に置き換えてください。
from dask.distributed import Client, LocalCluster
from dask_jobqueue import PBSCluster
from prefect import flow, task
from prefect_dask import DaskTaskRunner
cluster_kwargs = {
"cores": 1,
"memory": "1 GB",
"job_script_prologue" : ["source /etc/profile.d/modules.sh","module load python/3.11/3.11.14"],
"walltime": "00:05:00",
"n_workers" : 1,
"local_directory" : "~/tmp",
"resource_spec" : "rt_QF=1",
"job_extra_directives" : ["-W group_list=<group_name>"]
}
cluster = PBSCluster(**cluster_kwargs)
pbs_jobs = 2
cluster.scale(jobs=pbs_jobs)
client = Client(cluster)
def fib(n: int) -> int:
if n <= 2:
return 1
return fib(n-1) + fib(n-2)
@task
def task1(_: int) -> int:
return fib(5)
@task
def task_sum(ls: list[int]) -> int:
ret = 0
for x in ls:
ret += x
# タスクの停止確認のためにsleep処理を追記
import time
time.sleep(100)
return ret
@flow(task_runner=DaskTaskRunner(address=client.scheduler.address), log_prints=True)
def main():
xs = []
for i in range(3):
xs.append(task1.submit(i))
print(task_sum.submit(xs).result())
if __name__ == "__main__":
main.serve(name="my-pbs-example-deployment")
deploymentの実行
作成したワークフローを deployment として登録します。
(venv_name) [username@qh001 ~]$ python test_sysH_deployment.py &
[1] 2987156
(venv_name) [username@qh001 ~]$ Your flow 'main' is being served and polling for scheduled runs!
To trigger a run for this flow, use the following command:
$ prefect deployment run 'main/my-pbs-example-deployment'
You can also run your flow via the Prefect UI: http://10.0.0.33:4200/deployments/deployment/19211299-0bc5-4232-a732-83e3549e594f
(venv_name) [username@qh001 ~]$
ワークフローの実行
表示されたコマンド prefect deployment run 'main/my-pbs-example-deployment' を実行するとワークフローが実行されます
(venv_name) [username@qh001 ~]$ prefect deployment run 'main/my-pbs-example-deployment'
Creating flow run for deployment 'main/my-pbs-example-deployment'...
Created flow run 'dashing-quetzal'.
└── UUID: b8fb39c9-061a-4cf3-b9c3-ebd0f29dd3b7
└── Parameters: {}
└── Job Variables: {}
└── Scheduled start time: 2026-01-28 16:39:18 JST (now)
└── URL: http://10.0.0.33:4200/runs/flow-run/b8fb39c9-061a-4cf3-b9c3-ebd0f29dd3b7
(venv_name) [username@qh001 ~]$ 16:39:21.721 | INFO | prefect.flow_runs.runner - Runner 'my-pbs-example-deployment' submitting flow run 'b8fb39c9-061a-4cf3-b9c3-ebd0f29dd3b7'
16:39:21.777 | INFO | prefect.flow_runs.runner - Completed submission of flow run 'b8fb39c9-061a-4cf3-b9c3-ebd0f29dd3b7'
16:39:24.501 | INFO | Flow run 'dashing-quetzal' - Beginning flow run 'dashing-quetzal' for flow 'main'
16:39:24.504 | INFO | Flow run 'dashing-quetzal' - View at http://10.0.0.33:4200/runs/flow-run/b8fb39c9-061a-4cf3-b9c3-ebd0f29dd3b7
16:39:24.524 | INFO | prefect.task_runner.dask - Connecting to existing Dask cluster PBSCluster(383d0cbb, 'tcp://10.0.18.52:33533', workers=0, threads=0, memory=0 B)
(venv_name) [username@qh001 ~]$
ワークフローの停止
ワークフロー実行時に表示された UUID を指定して prefect flow-run cancel b8fb39c9-061a-4cf3-b9c3-ebd0f29dd3b7 を実行すると、ワークフローを停止することができます。
ここでは、コマンドを使ってワークフローを停止しましたが、Prefect UI からでもワークフローを停止することが可能です。
(venv_name) [username@qh001 ~]$ prefect flow-run cancel b8fb39c9-061a-4cf3-b9c3-ebd0f29dd3b7
Flow run 'b8fb39c9-061a-4cf3-b9c3-ebd0f29dd3b7' was successfully scheduled for cancellation.
16:40:33.376 | INFO | prefect.flow_runs.runner - Process for flow run 'dashing-quetzal' exited with status code: -15; This indicates that the process exited due to a SIGTERM signal. Typically, this is caused by manual cancellation.
16:40:36.388 | INFO | prefect.flow_runs.runner - Cancelled flow run 'dashing-quetzal'!
(venv_name) [username@qh001 ~]$
システム F でタスク実行
システム F での量子回路実行タスクの実行方法について説明します。
ワークフロープログラムの作成
システム F で量子回路実行タスクを実行するワークフロープログラムを作成します。
ワークフロープログラムとしては下記を行っています。
- システム H において量子回路 (
OpenQASM3形式) を作成 - システム H から REST API 経由で下記を実行
- システム F で量子回路実行、結果取得
- プログラムの終了時にシステム F の量子回路実行が終了していない場合は、実行をキャンセル
次のプログラムをtest_sysF_task.pyとして保存します。
プログラム内の<oqtopus_file>はご自身の環境の.oqtopusファイルのパスに置き換えてください。
from prefect import flow, task
from quri_parts_oqtopus.backend import (
OqtopusConfig,
OqtopusSamplingBackend,
OqtopusSamplingResult)
from quri_parts.backend import BackendError
import sys
@task
def get_circuit() ->str:
program = '''OPENQASM 3;
include "stdgates.inc";
qubit[2] q;
bit[2] c;
h q[0];
cx q[0], q[1];
c[0] = measure q[0];
c[1] = measure q[1];
'''
return program
@task
def do_system_f_job(circuit:str) -> OqtopusSamplingResult:
on_exec = False
try:
backend = OqtopusSamplingBackend(
OqtopusConfig.from_file("oqtopus-dev", path="<oqtopus_file>"))
transpiler_info = {
"initial_layout" : "[0,1]",
}
job = backend.sample_qasm(
circuit,
device_id = "abciq-f-1",
transpiler_info = transpiler_info,
shots=10000,
)
on_exec = True
jobid = job.job_id
print(f"job_id: {jobid}")
res = job.result()
on_exec = False
print("===> START: OqtopusSamplingJob")
print(job)
print("<=== END : OqtopusSamplingJob")
return res
except BackendError:
print(f"message: {job.job_info['message']}", file=sys.stderr)
raise
finally:
# Check system F job
if on_exec:
if (jobid is not None):
status = backend.retrieve_job(jobid).status
isFinish = ( status == "succeeded" or status == "failed")
if not isFinish:
print(f"job status@finally: {status}")
print("===> START: job cancel")
job.cancel()
print("<=== END : job cancel")
print(f"job status: {job.status}")
@flow(log_prints=True)
def main():
circuit = get_circuit()
res = do_system_f_job(circuit)
print(res)
if __name__ == "__main__":
main()
ワークフローの実行
作成したワークフロープログラムは、以下のように実行します。
(venv_name) [username@qh001 test]$ python test_sysF_task.py
また、実行したワークフローは Prefect UI から確認することができます。
ワークフローの停止
Prefect のワークフローを停止するためには deployment という仕組みを利用する必要があります。
ワークフロープログラムの作成
システム F で量子回路実行タスクを停止できることを確認するためのワークフロープログラムを作成します。フロー main() に対して serve というメソッドを呼び出すことで deployment を作成するようにしています。
次のプログラムをtest_sysF_deployment.pyとして保存します。
プログラム内の<oqtopus_file>はご自身の環境の.oqtopusファイルのパスに置き換えてください。
from prefect import flow, task
from quri_parts_oqtopus.backend import (
OqtopusConfig,
OqtopusSamplingBackend,
OqtopusSamplingResult)
from quri_parts.backend import BackendError
import sys
@task
def get_circuit() ->str:
program = '''OPENQASM 3;
include "stdgates.inc";
qubit[2] q;
bit[2] c;
h q[0];
cx q[0], q[1];
c[0] = measure q[0];
c[1] = measure q[1];
'''
return program
@task
def do_system_f_job(circuit:str) -> OqtopusSamplingResult:
on_exec = False
try:
backend = OqtopusSamplingBackend(
OqtopusConfig.from_file("oqtopus-dev", path="<oqtopus_file>"))
transpiler_info = {
"initial_layout" : "[0,1]",
}
# 量子回路の実行
job = backend.sample_qasm(
circuit,
device_id = "abciq-f-1",
transpiler_info = transpiler_info,
shots=10000,
)
on_exec = True
jobid = job.job_id
print(f"job_id: {jobid}")
# タスクの停止確認のためにsleep処理を追記
import time
for i in range(10):
time.sleep(2)
print(f"sleep processing... {i+1}/10")
res = job.result()
on_exec = False
print("===> START: OqtopusSamplingJob")
print(job)
print("<=== END : OqtopusSamplingJob")
return res
except BackendError:
print(f"message: {job.job_info['message']}", file=sys.stderr)
raise
finally:
# Check system F job
if on_exec:
if (jobid is not None):
status = backend.retrieve_job(jobid).status
isFinish = ( status == "succeeded" or status == "failed")
if not isFinish:
print(f"job status@finally: {status}") # DEBUG
print("===> START: job cancel") # DEBUG
job.cancel()
print("<=== END : job cancel") # DEBUG
print(f"job status: {job.status}")
@flow(log_prints=True)
def main():
circuit = get_circuit()
res = do_system_f_job(circuit)
print(res)
if __name__ == "__main__":
main.serve(name="my-system-f-job-deployment")
deploymentの実行
作成したワークフローを deployment として登録します。
(venv_name) [username@qh001 ~]$ python test_sysF_deployment.py &
[1] 3008604
(venv_name) [username@qh001 ~]$ Your flow 'main' is being served and polling for scheduled runs!
To trigger a run for this flow, use the following command:
$ prefect deployment run 'main/my-system-f-job-deployment'
You can also run your flow via the Prefect UI: http://10.0.0.33:4200/deployments/deployment/a2476573-575c-4367-b7ff-58d61206c543
(venv_name) [username@qh001 ~]$
ワークフローの実行
表示されたコマンド prefect deployment run 'main/my-system-f-job-deployment' を実行するとワークフローが実行されます。
(venv_name) [username@qh001 ~]$ prefect deployment run 'main/my-system-f-job-deployment'
Creating flow run for deployment 'main/my-system-f-job-deployment'...
Created flow run 'fractal-falcon'.
└── UUID: c902fe93-bcd8-42da-89ab-74da93015dd8
└── Parameters: {}
└── Job Variables: {}
└── Scheduled start time: 2026-01-28 17:17:02 JST (now)
└── URL: http://10.0.0.33:4200/runs/flow-run/c902fe93-bcd8-42da-89ab-74da93015dd8
ワークフローの停止
ワークフロー実行時に表示された UUID を指定して prefect flow-run cancel c902fe93-bcd8-42da-89ab-74da93015dd8 を実行すると、ワークフローを停止することができます。
ここでは、コマンドを使ってワークフローを停止しましたが、Prefect UI からでもワークフローを停止することが可能です。
(venv_name) [username@qh001 ~]$ prefect flow-run cancel c902fe93-bcd8-42da-89ab-74da93015dd8
Flow run 'c902fe93-bcd8-42da-89ab-74da93015dd8' was successfully scheduled for cancellation.
17:17:16.596 | ERROR | Task run 'do_system_f_job-a29' - Crash detected! Execution was aborted by a termination signal.
17:17:16.597 | ERROR | Task run 'do_system_f_job-a29' - Finished in state Crashed('Execution was aborted by a termination signal.')
17:17:16.609 | INFO | prefect.flow_runs.runner - Process for flow run 'fractal-falcon' exited with status code: -15; This indicates that the process exited due to a SIGTERM signal. Typically, this is caused by manual cancellation.
(venv_name) [username@qh001 ~]$ 17:17:19.606 | INFO | prefect.flow_runs.runner - Cancelled flow run 'fractal-falcon'!
(venv_name) [username@qh001 ~]$
Singularity コンテナを用いたタスク実行
prefect-singularity では @singularity_task デコレーターによって Singularity コンテナ内でタスク実行できます。
利用したい Singularity イメージを指定することで、その環境でタスク実行ができます。
Caution
Singularity イメージには Python3.11, cloudpickle 3.1.2 がインストールされている必要があります。 Python のバージョンや cloudpickle のバージョンが異なると、動作しない可能性があります。
設定
@singularity_task デコレーターによって Singularity コンテナ内でタスク実行するためには、設定用のクラスSingularityConfigをインスタンス化し各種設定を行います。
設定用のクラスSingularityConfigの設定項目は以下の通りです。
| 大項目 | 小項目 | 説明 | 設定例 | デフォルト値 |
|---|---|---|---|---|
| executor | PBS を介した実行か、ローカルでの実行か | local または pbs | local | |
| mpi_setting | MPI実行用の設定 | |||
| mpi_type | MPIの種別 | "openmpi","hpc-x","intel-mpi" | None | |
| mpi_ops | MPIのオプション | "-n 8 --map-by ppr:4:node --display map" | None | |
| pbsjob_setting | PBS を介したジョブ実行用の設定 | |||
| prerun_command | タスクを実行前に実行するコマンド | ["module load ABC", "source ~/venv/bin/activate"] | None | |
| qsub_args | qsub 実行時の引数 | [" -W group_list=group", "-l resource_type=num:mpiprocs=num"] | None | |
| singularity_setting | singularity の設定 | |||
| singularity_image_path | singularity イメージの Path | "~/path/to/singularity_image_path" | None | |
| gpu | GPU 実行か否か | True | None | |
| prefect_kwargs | Prefect のタスクの設定値1 | { "retries" : 3 } |
None |
Caution
- 現在、
@singularity_taskデコレーターは PBS に対応しておりません。そのため、executorを local 以外に設定するとエラーとなります。また、pbsjob_settingの設定値は無視されます。 mpi_settingやmpi_typeが設定されていない場合は、逐次実行として扱います。- 現在、
mpi_typeはOpen MPIのみサポートします。
ワークフローの実行
cuQuantum Appliance を Singularity イメージとして指定して、 Qiskit Aer による量子シミュレーションタスクを実行する方法について説明します。
Singularityイメージの作成
NVIDIA NGC で公開されている cuQuantum Appliance の Docker イメージから Singularity イメージを作成します。 (現在サポートしているのは、cuQuantum appliance version 25.11 です。)
Singularity イメージ作成手順はcuQuantum Applianceを参照してください。
コンテナイメージを作成するためのDefinition fileは以下のものを使用してください。
[username@qes01 ~]$ cat qc.def
Bootstrap: docker
From: nvcr.io/nvidia/cuquantum-appliance:25.11-cuda12.9.1-devel-ubuntu24.04-x86_64
%environment
source /opt/etc/bashrc
export CONDA_ENV_NAME=cuquantum
conda info -e >> /dev/stderr
%post -c /bin/bash
apt update
apt dist-upgrade -y
apt install git vim -y
chmod -R o+rX /home/cuquantum
mkdir -p /opt/etc
echo -e "#! /bin/bash\n\n# script to activate the conda environment" > ~/.bashrc \
&& conda init bash \
&& echo -e "\nconda activate cuquantum" >> ~/.bashrc \
&& echo "echo \"Hello World\" >>/dev/stderr " >> ~/.bashrc \
&& conda clean -a \
&& cp ~/.bashrc /opt/etc/bashrc
/bin/bash -rcfile /opt/etc/bashrc -c "pip install cloudpickle==3.1.2"
標準のDefinition fileとの差分は、
/bin/bash -rcfile /opt/etc/bashrc -c "pip install cloudpickle==3.1.2"
のみです。
Qiskit Aer による量子回路シミュレーションタスクの実行
cuQuantum Appliance の Qiskit Aer で 20 量子ビットの GHZ 回路の量子シミュレーションタスクを実行するワークフロープログラムの例を示します。
ワークフロープログラムとしては、設定用のクラス SingularityConfig をインスタンス化し、Qiskit Aer による量子回路シミュレーションを実行する関数 ghz_circuit に@singularity_task デコレーターを付与してます。このとき、インスタンス化された SingularityConfig はデコレーターの引数として渡しています。
Singularityコンテナを用いたタスク実行のための設定としては、ローカル実行でシングル GPU を使った実行となるように設定を行っています。
また、量子回路シミュレーションを実行する関数 ghz_circuit は、シミュレーション結果を辞書型で return するようにしています。
次のプログラムをtest_singularity_task.pyとして保存します。
プログラム内の/path/to/singularity_imageはご自身の環境のqc.sifファイルのパスに置き換えてください。
from prefect_singularity import singularity_task, SingularityConfig
from prefect import flow
from qiskit import QuantumCircuit, transpile
from qiskit_aer import Aer
from qiskit.result import Result
config = SingularityConfig(
executor="local",
singularity_setting = {
"singularity_image_path" : "/path/to/singularity_image",
"gpu" : True
},
prefect_kwargs = {
"name" : "ghz_simulation"
}
)
@singularity_task(config)
def ghz_circuit(n_qubits):
circuit = QuantumCircuit(n_qubits)
circuit.h(0)
for qubit in range(n_qubits - 1):
circuit.cx(qubit, qubit + 1)
circuit.measure_all()
simulator = Aer.get_backend('aer_simulator_statevector')
circuit = transpile(circuit, simulator)
job = simulator.run(circuit)
result = job.result()
return result.to_dict()
@flow
def main():
res_data = ghz_circuit(20)
res = Result.from_dict(res_data)
print(res.get_counts())
if __name__ == "__main__":
main()
Caution
@singularity_task デコレーターで cuQuantum Appliance の Qiskit で量子回路シミュレーションを実行し、
その結果を return するタスク定義時には例のように辞書型で return するようにタスクを定義してください。
作成したワークフロープログラムを実行すると、以下のようになります。
想定される 00000000000000000000 と 11111111111111111111が同じくらいの頻度で出現する結果が得られていることがわかります。
また、実行したワークフローは Prefect UI から確認することができます。
(venv_name) [username@qh001 ~]$ python test_singularity_task.py
17:06:01.414 | INFO | Flow run 'ambitious-coati' - Beginning flow run 'ambitious-coati' for flow 'main'
17:06:01.418 | INFO | Flow run 'ambitious-coati' - View at http://10.0.0.33:4500/runs/flow-run/e520f69a-678d-4786-8cf1-9dba8973149f
17:06:23.892 | INFO | Task run 'ghz_simulation-cb1' - INFO: Setting 'NVIDIA_VISIBLE_DEVICES=all' to emulate legacy GPU binding.
INFO: Setting --writable-tmpfs (required by nvidia-container-cli)
Hello World
# conda environments:
#
# * -> active
# + -> frozen
base /opt/conda
cuquantum * /opt/conda/envs/cuquantum
17:06:23.901 | INFO | Task run 'ghz_simulation-cb1' - Finished in state Completed()
{'00000000000000000000': 500, '11111111111111111111': 524}
17:06:23.920 | INFO | Flow run 'ambitious-coati' - Finished in state Completed()
MPIを使ったマルチGPUによる量子回路シミュレーションタスクの実行
MPIを使用して、4 GPUで20 量子ビットの GHZ 回路の量子シミュレーションタスクを実行するワークフロープログラムの例を示します。
シングルGPUを使用した例との差分は、下記のconfig = SigularityConfig(内の設定のみです。
mpi_setting = {
"mpi_type" : "Openmpi",
"mpi_ops" : "-np 4 -display map"
},
次のプログラムをtest_singularity_mpi_task.pyとして保存します。
プログラム内の/path/to/singularity_imageはご自身の環境のqc.sifファイルのパスに置き換えてください。
from prefect_singularity import singularity_task, SingularityConfig
from prefect import flow
from qiskit import QuantumCircuit, transpile
from qiskit_aer import Aer
from qiskit.result import Result
config = SingularityConfig(
executor="local",
mpi_setting = {
"mpi_type" : "Openmpi",
"mpi_ops" : "-np 4 -display map"
},
singularity_setting = {
"singularity_image_path" : "/path/to/singularity_image",
"gpu" : True
},
prefect_kwargs = {
"name" : "ghz_simulation"
}
)
@singularity_task(config)
def ghz_circuit(n_qubits):
circuit = QuantumCircuit(n_qubits)
circuit.h(0)
for qubit in range(n_qubits - 1):
circuit.cx(qubit, qubit + 1)
circuit.measure_all()
simulator = Aer.get_backend('aer_simulator_statevector')
circuit = transpile(circuit, simulator)
job = simulator.run(circuit)
result = job.result()
return result.to_dict()
@flow
def main():
res_data = ghz_circuit(20)
res = Result.from_dict(res_data)
print(res.get_counts())
if __name__ == "__main__":
main()
MPIで複数プロセスを実行するためには、インタラクティブジョブ実行時にmpiprocs=指定を追加する必要があります。
[username@qes01 ~]$ qsub -l rt_QF=1:mpiprocs=4 -W group_list=grpname -I
ワークフロープログラムを実行する前に、openmpiモジュールのロードを行います。
(venv_name) [username@qh001 ~]$ module load openmpi/5.0.8
(venv_name) [username@qh001 ~]$ python test_singularity_mpi_task.py
ワークフローの停止
Prefect のワークフローを停止するためには deployment という仕組みを利用する必要があります。
ワークフロープログラムの作成
cuQuantum Appliance の Qiskit Aer で 20 量子ビットの GHZ 回路の量子シミュレーションタスクを停止できることを確認するためのワークフロープログラムを作成します。
フロー main() に対して serve というメソッドを呼び出すことで deployment を作成するようにしています。
次のプログラムをtest_singularity_deployment.pyとして保存します。
プログラム内の/path/to/singularity_imageはご自身の環境のqc.sifファイルのパスに置き換えてください。
from prefect_singularity import singularity_task, SingularityConfig
from prefect import flow
from qiskit import QuantumCircuit, transpile
from qiskit_aer import Aer
from qiskit.result import Result
config = SingularityConfig(
executor="local",
singularity_setting = {
"singularity_image_path" : "/path/to/singularity_image",
"gpu" : True
},
prefect_kwargs = {
"name" : "ghz_simulation"
}
)
@singularity_task(config)
def ghz_circuit(n_qubits):
circuit = QuantumCircuit(n_qubits)
circuit.h(0)
for qubit in range(n_qubits - 1):
circuit.cx(qubit, qubit + 1)
circuit.measure_all()
simulator = Aer.get_backend('aer_simulator_statevector')
circuit = transpile(circuit, simulator)
job = simulator.run(circuit)
# タスクの停止確認のためにsleep処理を追記
import time
time.sleep(100)
result = job.result()
return result.to_dict()
@flow
def main():
res_data = ghz_circuit(20)
res = Result.from_dict(res_data)
print(res.get_counts())
if __name__ == "__main__":
main.serve(name="ghz-simulation-deployment")
deploymentの実行
作成したワークフローを deployment として登録します。
(venv_name) [username@qh001 ~]$ python test_singularity_deployment.py &
[1] 1764126
(venv_name) [username@qh001 ~]$ Your flow 'main' is being served and polling for scheduled runs!
To trigger a run for this flow, use the following command:
$ prefect deployment run 'main/ghz-simulation-deployment'
You can also run your flow via the Prefect UI: http://10.0.0.33:4500/deployments/deployment/5b2c68e4-53f1-4885-a019-2dbd2c8af2b8
ワークフローの実行
表示されたコマンド prefect deployment run 'main/ghz-simulation-deployment' を実行するとワークフローが実行されます。
(venv_name) [username@qh001 ~]$ prefect deployment run 'main/ghz-simulation-deployment'
Creating flow run for deployment 'main/ghz-simulation-deployment'...
Created flow run 'emerald-cricket'.
└── UUID: d1f8f229-d06c-4dab-bc43-6cec010122ac
└── Parameters: {}
└── Job Variables: {}
└── Scheduled start time: 2026-01-27 17:15:36 JST (now)
└── URL: http://10.0.0.33:4500/runs/flow-run/d1f8f229-d06c-4dab-bc43-6cec010122ac
ワークフローの停止
ワークフロー実行時に表示された UUID を指定して prefect flow-run cancel d1f8f229-d06c-4dab-bc43-6cec010122ac を実行すると、ワークフローを停止することができます。
ここでは、コマンドを使ってワークフローを停止しましたが、Prefect UI からでもワークフローを停止することが可能です。
(venv_name) [username@qh001 ~]$ prefect flow-run cancel d1f8f229-d06c-4dab-bc43-6cec010122ac
Flow run 'd1f8f229-d06c-4dab-bc43-6cec010122ac' was successfully scheduled for cancellation.
17:16:34.292 | ERROR | Task run 'ghz_simulation-421' - Crash detected! Execution was aborted by a termination signal.
17:16:34.293 | ERROR | Task run 'ghz_simulation-421' - Finished in state Crashed('Execution was aborted by a termination signal.')
17:16:34.300 | INFO | prefect.flow_runs.runner - Process for flow run 'emerald-cricket' exited with status code: -15; This indicates that the process exited due to a SIGTERM signal. Typically, this is caused by manual cancellation.
(venv_name) [username@qh001 ~]$ 17:16:37.319 | INFO | prefect.flow_runs.runner - Cancelled flow run 'emerald-cricket'!
-
詳細は Prefect のドキュメントの Task configuration を参照してください ↩