このガイドでは Python Workflows SDK を取り上げ、Python でワークフローを組み立てて作成する方法を説明します。
WorkflowEntrypoint は、Python ワークフローのメインエントリポイントです。WorkflowEntrypoint クラスを拡張し、run メソッドを実装します。
from workers import WorkflowEntrypoint
class MyWorkflow(WorkflowEntrypoint):
async def run(self, event, step):
# steps here-
step.do(name=None, *, concurrent=False, config=None)— ワークフロー内のステップを定義するデコレーターです。name— ステップの任意の名前です。省略すると、関数名(func.__name__)が使われます。concurrent— このステップの依存関係を同時実行できるかどうかを示す、任意のブール値です。config— ステップ固有のリトライ動作 を設定する、任意のWorkflowStepConfigです。Python の辞書として渡し、その後WorkflowStepConfigオブジェクトへ型変換されます。
name以外のパラメーターは、キーワード専用です。
依存関係は、パラメーター名から暗黙に解決されます。ステップ関数のパラメーター名が、以前に宣言したステップ関数と一致する場合、その結果がステップへ注入されます。
ctx パラメーターを定義すると、その引数にステップコンテキストが注入されます。
from workers import WorkflowEntrypoint
class MyWorkflow(WorkflowEntrypoint):
async def run(self, event, step):
@step.do()
async def my_first_step():
# do some work
return "Hello World!"
await my_first_step()デコレーターはステップを呼び出すのではなく、ステップを起動できる callable を返すだけです。ステップを実行するには、その callable を呼ぶ必要があります。
ステップから状態を返すときは、戻り値がシリアライズ可能である必要があります。
-
step.sleep(name, duration)name— ステップの名前です。duration— スリープする期間です。ミリ秒のnumber、またはWorkflowDuration互換の文字列です。
async def run(self, event, step):
await step.sleep("my-sleep-step", "10 seconds")-
step.sleep_until(name, timestamp)name— ステップの名前です。timestamp— ワークフローインスタンスをスリープさせる先のdatetime.datetimeオブジェクト、または Unix エポックからの秒数です。
import datetime
async def run(self, event, step):
await step.sleep_until("my-sleep-step", datetime.datetime.now() + datetime.timedelta(seconds=10))-
step.wait_for_event(name, event_type, timeout="24 hours")name— ステップの名前です。event_type— 待機するイベントの種類です。timeout—wait_for_event呼び出しのタイムアウトです。既定のタイムアウトは 24 時間です。
async def run(self, event, step):
await step.wait_for_event("my-wait-for-event-step", "my-event-type")event パラメーターは、ワークフローインスタンスへ渡されたペイロードと、そのほかのメタデータを含む辞書です。
payload- ワークフローインスタンスへ渡されたペイロードです。timestamp- ワークフローがトリガーされた時刻です。instanceId- 現在のワークフローインスタンスの ID です。workflowName- ワークフローの名前です。
Workflows のセマンティクスでは、トップレベルまで伝播した例外を捕捉できます。
except ブロックで特定の例外だけを捕捉しても、うまくいかないことがあります。一部の Python エラーは、RPC 層を通過するときに同じ種類のエラーとして再インスタンス化されないためです。
async def run(self, event, step):
async def try_step(fn):
try:
return await fn()
except Exception as e:
print(f"Successfully caught {type(e).__name__}: {e}")
@step.do("my_failing")
async def my_failing():
print("Executing my_failing")
raise TypeError("Intentional error in my_failing")
await try_step(my_failing)Python Workflows SDK は、ステップをリトライしないことを示す NonRetryableError クラスを提供します。
from workers.workflows import NonRetryableError
raise NonRetryableError(message)step.do デコレーターの config パラメーターに WorkflowStepConfig オブジェクトを渡すと、ステップを特定のリトライポリシーにバインドできます。
Python Workflows では、dict が WorkflowStepConfig 型に従っている必要があります。
from workers import WorkflowEntrypoint
class DemoWorkflowClass(WorkflowEntrypoint):
async def run(self, event, step):
@step.do('step-name', config={"retries": {"limit": 1, "delay": "10 seconds"}})
async def first_step():
# do some work
passctx パラメーターを定義すると、その引数に ステップコンテキスト が注入されます。コンテキストは、次のキーを持つ辞書です。
| キー | 型 | 説明 |
|---|---|---|
step |
dict |
name(ステップ名)と count(この名前で step.do が呼ばれた回数)を含みます。 |
attempt |
int |
現在の試行回数です(1 始まり)。 |
config |
dict |
このステップの解決済みリトライおよびタイムアウト設定です。 |
from workers import WorkflowEntrypoint
class CtxWorkflow(WorkflowEntrypoint):
async def run(self, event, step):
@step.do()
async def read_context(ctx):
print(ctx["step"]["name"]) # step name
print(ctx["step"]["count"]) # step count
print(ctx["attempt"]) # attempt number
print(ctx["config"]) # resolved step config
return ctx["attempt"]
return await read_context()env は、JsProxy ↗ 経由で Python スクリプトに公開される JavaScript オブジェクトです。JavaScript の Worker と同じようにバインディングへアクセスできます。利用できるメソッドは、Workflow バインディングのドキュメント を参照してください。
さきほどのバインディング名を MY_WORKFLOW とします。新しいインスタンスの作成方法は次のとおりです。
from workers import Response, WorkerEntrypoint
class Default(WorkerEntrypoint):
async def fetch(self, request):
instance = await self.env.MY_WORKFLOW.create()
return Response.json({"status": "success"})