コンテンツにスキップ

ワークフローと実行例

ワークフローの操作について

ここでは、ワークフローを以下の方法で操作します。

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_kwargsdask_jobqueue.PBSCluster の設定を記入
  • dask_jobqueue.PBSCluster のインスタンス作成時の引数として cluster_kwargsを指定
  • dask_jobqueue.PBSClusterscaleメソッドの引数で、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_settingmpi_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 するようにタスクを定義してください。

作成したワークフロープログラムを実行すると、以下のようになります。 想定される 0000000000000000000011111111111111111111が同じくらいの頻度で出現する結果が得られていることがわかります。

また、実行したワークフローは 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'!

  1. 詳細は Prefect のドキュメントの Task configuration を参照してください