サイトアイコンmaita tomoya dev io

AWS Step Functions

DevOps/インフラ

AWS Step Functions

ワークフローオーケストレーションとは

ワークフローオーケストレーションとは、複数の処理ステップを定義された順序で実行し、ステップ間のデータ受け渡し、エラー処理、リトライなどを管理する仕組み。オーケストラの指揮者(オーケストレーター)が各楽器の演奏タイミングを制御するように、複数のサービスの実行を調整する。

なぜオーケストレーションが必要か

単一のLambda関数で全てを処理すると以下の問題が発生する。

問題説明
タイムアウトLambda15分制限では長い処理が完了しない
エラー処理の複雑化try-catchのネストが深くなり可読性が低下
リトライの管理どのステップまで成功したか追跡が困難
状態の管理中間結果の受け渡しが煩雑になる
可視性の欠如処理のどこで止まったか把握しにくい

Step Functionsはこれらの課題を解決する。


AWS Step Functionsとは

AWS Step Functionsは、AWSが提供するサーバーレスのワークフローオーケストレーションサービス。視覚的なワークフローを使って、複数のAWSサービスを組み合わせた分散アプリケーションを構築できる。

ステートマシンと呼ばれるワークフローを定義し、各ステート(状態)がLambda関数の呼び出し、DynamoDBへの書き込み、SNS通知などの処理を実行する。

Step Functionsの全体像

graph TD A[開始] --> B{注文内容の検証} B -->|有効| C[在庫確認] B -->|無効| G[エラー通知] C -->|在庫あり| D[決済処理] C -->|在庫なし| G D -->|成功| E[配送手配] D -->|失敗| F[決済エラー処理] E --> H[完了通知] F --> G G --> I[終了] H --> I

このようなワークフローをJSON(ASL: Amazon States Language)で定義し、Step Functionsが自動的にステップ間の遷移、エラー処理、リトライを管理する。


Standard vs Express ワークフロー

Step Functionsには2つのワークフロータイプがある。

比較表

項目StandardExpress
最大実行時間1年5分
実行モデル正確に1回(Exactly-once)最低1回(At-least-once)または最大1回
実行履歴Step Functions コンソールで確認可能CloudWatch Logsに出力
料金モデル状態遷移ごとに課金実行回数×実行時間で課金
料金$0.025 / 1,000状態遷移$1.00 / 100万リクエスト + 実行時間
実行開始レート2,000/秒100,000/秒以上
適したユースケース長時間・低頻度の処理短時間・高頻度の処理

Express ワークフローのサブタイプ

タイプ説明ユースケース
同期(Synchronous)呼び出し元がレスポンスを待つAPI Gatewayのバックエンド
非同期(Asynchronous)即座に実行IDを返すイベント駆動処理

選択の指針

graph TD A[ワークフロータイプの選択] --> B{実行時間は5分以内?} B -->|いいえ| C[Standard] B -->|はい| D{高スループットが必要?} D -->|はい| E[Express] D -->|いいえ| F{Exactly-once が必要?} F -->|はい| C F -->|いいえ| G{コスト重視?} G -->|状態遷移が多い| E G -->|状態遷移が少ない| C

ステートの種類

Step Functionsのワークフローは、様々な種類のステート(状態)で構成される。

全ステート一覧

ステート説明用途
Task処理を実行するLambda呼び出し、AWS API実行
Choice条件分岐if-else的なロジック
Parallel並列実行複数の処理を同時実行
Map反復処理配列の各要素を並列処理
Wait待機一定時間待つ、特定時刻まで待つ
Passパススルーデータの変換、デバッグ
Succeed成功終了ワークフローの正常終了
Fail失敗終了ワークフローの異常終了

Task ステート

最も重要なステート。AWSサービスの呼び出しを行う。

{
  "Type": "Task",
  "Resource": "arn:aws:states:::lambda:invoke",
  "Parameters": {
    "FunctionName": "arn:aws:lambda:ap-northeast-1:123456789:function:ProcessOrder",
    "Payload.$": "$"
  },
  "ResultPath": "$.orderResult",
  "Next": "CheckResult"
}

SDK統合(直接AWS APIを呼び出す)

Lambda関数を経由せずにAWSサービスを直接呼び出せる。200以上のAWSサービスに対応。

{
  "Type": "Task",
  "Resource": "arn:aws:states:::dynamodb:putItem",
  "Parameters": {
    "TableName": "Orders",
    "Item": {
      "OrderId": { "S.$": "$.orderId" },
      "Status": { "S": "CREATED" },
      "CreatedAt": { "S.$": "$$.State.EnteredTime" }
    }
  },
  "Next": "NotifyUser"
}

Choice ステート

{
  "Type": "Choice",
  "Choices": [
    {
      "Variable": "$.orderTotal",
      "NumericGreaterThan": 10000,
      "Next": "ApplyDiscount"
    },
    {
      "Variable": "$.memberType",
      "StringEquals": "PREMIUM",
      "Next": "PremiumProcessing"
    }
  ],
  "Default": "StandardProcessing"
}

Parallel ステート

{
  "Type": "Parallel",
  "Branches": [
    {
      "StartAt": "SendEmail",
      "States": {
        "SendEmail": {
          "Type": "Task",
          "Resource": "arn:aws:states:::sns:publish",
          "Parameters": {
            "TopicArn": "arn:aws:sns:ap-northeast-1:123456789:OrderNotification",
            "Message.$": "$.message"
          },
          "End": true
        }
      }
    },
    {
      "StartAt": "UpdateInventory",
      "States": {
        "UpdateInventory": {
          "Type": "Task",
          "Resource": "arn:aws:states:::lambda:invoke",
          "Parameters": {
            "FunctionName": "UpdateInventory",
            "Payload.$": "$"
          },
          "End": true
        }
      }
    }
  ],
  "Next": "OrderComplete"
}

Map ステート

配列の各要素に対して同じ処理を実行する。Inline Map(小規模)とDistributed Map(大規模)がある。

{
  "Type": "Map",
  "ItemsPath": "$.orderItems",
  "ItemProcessor": {
    "ProcessorConfig": {
      "Mode": "INLINE"
    },
    "StartAt": "ProcessItem",
    "States": {
      "ProcessItem": {
        "Type": "Task",
        "Resource": "arn:aws:states:::lambda:invoke",
        "Parameters": {
          "FunctionName": "ProcessItem",
          "Payload.$": "$"
        },
        "End": true
      }
    }
  },
  "MaxConcurrency": 10,
  "Next": "AggregateResults"
}

Distributed Mapは、S3バケット内の数百万オブジェクトを対象とした大規模並列処理に対応する。最大10,000の並列実行が可能。


Amazon States Language(ASL)

ASLは、Step Functionsのワークフローを定義するためのJSONベースの仕様。

完全なワークフロー例

{
  "Comment": "注文処理ワークフロー",
  "StartAt": "ValidateOrder",
  "States": {
    "ValidateOrder": {
      "Type": "Task",
      "Resource": "arn:aws:states:::lambda:invoke",
      "Parameters": {
        "FunctionName": "ValidateOrder",
        "Payload.$": "$"
      },
      "ResultPath": "$.validation",
      "Next": "IsOrderValid",
      "Catch": [
        {
          "ErrorEquals": ["States.ALL"],
          "Next": "OrderFailed",
          "ResultPath": "$.error"
        }
      ]
    },
    "IsOrderValid": {
      "Type": "Choice",
      "Choices": [
        {
          "Variable": "$.validation.Payload.isValid",
          "BooleanEquals": true,
          "Next": "ProcessPayment"
        }
      ],
      "Default": "OrderFailed"
    },
    "ProcessPayment": {
      "Type": "Task",
      "Resource": "arn:aws:states:::lambda:invoke",
      "Parameters": {
        "FunctionName": "ProcessPayment",
        "Payload.$": "$"
      },
      "ResultPath": "$.payment",
      "Retry": [
        {
          "ErrorEquals": ["States.TaskFailed"],
          "IntervalSeconds": 3,
          "MaxAttempts": 3,
          "BackoffRate": 2.0
        }
      ],
      "Next": "FulfillOrder",
      "Catch": [
        {
          "ErrorEquals": ["States.ALL"],
          "Next": "RefundPayment",
          "ResultPath": "$.error"
        }
      ]
    },
    "FulfillOrder": {
      "Type": "Task",
      "Resource": "arn:aws:states:::lambda:invoke",
      "Parameters": {
        "FunctionName": "FulfillOrder",
        "Payload.$": "$"
      },
      "Next": "OrderSucceeded"
    },
    "RefundPayment": {
      "Type": "Task",
      "Resource": "arn:aws:states:::lambda:invoke",
      "Parameters": {
        "FunctionName": "RefundPayment",
        "Payload.$": "$"
      },
      "Next": "OrderFailed"
    },
    "OrderSucceeded": {
      "Type": "Succeed"
    },
    "OrderFailed": {
      "Type": "Fail",
      "Error": "OrderProcessingError",
      "Cause": "Order processing failed"
    }
  }
}

入出力処理

ASLでは各ステートの入出力を制御するフィールドがある。

フィールド説明
InputPathステートに渡す入力をフィルタ"$.order"
Parameters入力を加工してタスクに渡すキーを組み替える
ResultSelectorタスク結果をフィルタ必要なフィールドだけ抽出
ResultPath結果を元の入力のどこに格納するか"$.taskResult"
OutputPath次のステートに渡す出力をフィルタ"$.taskResult"

データフローの順序:

入力 → InputPath → Parameters → [Task実行] → ResultSelector → ResultPath → OutputPath → 出力

エラー処理

Retry(リトライ)

"Retry": [
  {
    "ErrorEquals": ["CustomTransientError"],
    "IntervalSeconds": 1,
    "MaxAttempts": 5,
    "BackoffRate": 2.0,
    "JitterStrategy": "FULL"
  },
  {
    "ErrorEquals": ["States.ALL"],
    "IntervalSeconds": 5,
    "MaxAttempts": 2,
    "BackoffRate": 1.5
  }
]
パラメータ説明
ErrorEqualsリトライ対象のエラー名リスト
IntervalSeconds初回リトライまでの待機秒数
MaxAttempts最大リトライ回数(デフォルト3)
BackoffRateリトライ間隔の増加率(デフォルト2.0)
JitterStrategyジッターの適用戦略(FULL or NONE)

Catch(キャッチ)

リトライ後も失敗した場合のフォールバック先を定義する。

"Catch": [
  {
    "ErrorEquals": ["PaymentDeclined"],
    "Next": "HandleDeclinedPayment",
    "ResultPath": "$.error"
  },
  {
    "ErrorEquals": ["States.ALL"],
    "Next": "HandleUnknownError",
    "ResultPath": "$.error"
  }
]

定義済みエラー

エラー名説明
States.ALLすべてのエラーにマッチ
States.Timeoutタイムアウト
States.TaskFailedタスク実行の失敗
States.Permissions権限不足
States.ResultPathMatchFailureResultPathの適用失敗
States.ParameterPathFailureParametersのパス解決失敗
States.BranchFailedParallel/Mapのブランチ失敗
States.NoChoiceMatchedChoiceでマッチする条件なし
States.IntrinsicFailure組み込み関数の失敗

料金体系

Standard ワークフロー

料金 = 状態遷移数 × $0.025 / 1,000遷移
無料枠: 月4,000遷移

Express ワークフロー

料金 = リクエスト料金 + 実行時間料金
リクエスト料金 = $1.00 / 100万リクエスト
実行時間料金 = メモリ使用量 × 実行時間 × 単価

料金例

シナリオタイプ状態遷移/月月額料金(概算)
小規模注文処理Standard10万遷移約$2.50
中規模データパイプラインStandard100万遷移約$25
高頻度API処理Express1,000万リクエスト約$10 + 実行時間

実践的な設計パターン

Sagaパターン(補償トランザクション)

分散トランザクションにおいて、途中でエラーが発生した場合に、それまでに実行した処理を逆順で取り消す(補償する)パターン。

graph TD A[予約作成] --> B[決済処理] B --> C[在庫確保] C --> D[配送手配] B -->|失敗| E[予約キャンセル] C -->|失敗| F[決済返金] F --> E D -->|失敗| G[在庫解放] G --> F

人間承認パターン

ワークフローの途中で人間の承認を待つパターン。タスクトークンを使用する。

{
  "WaitForApproval": {
    "Type": "Task",
    "Resource": "arn:aws:states:::lambda:invoke.waitForTaskToken",
    "Parameters": {
      "FunctionName": "SendApprovalRequest",
      "Payload": {
        "taskToken.$": "$$.Task.Token",
        "orderId.$": "$.orderId"
      }
    },
    "TimeoutSeconds": 86400,
    "Next": "ProcessApprovedOrder"
  }
}

承認者がAPI経由でSendTaskSuccessまたはSendTaskFailureを呼び出すと、ワークフローが再開される。

動的並列処理パターン

Distributed Mapを使ったS3内の大量ファイル処理。

{
  "ProcessFiles": {
    "Type": "Map",
    "ItemProcessor": {
      "ProcessorConfig": {
        "Mode": "DISTRIBUTED",
        "ExecutionType": "STANDARD"
      },
      "StartAt": "TransformFile",
      "States": {
        "TransformFile": {
          "Type": "Task",
          "Resource": "arn:aws:states:::lambda:invoke",
          "Parameters": {
            "FunctionName": "TransformFile",
            "Payload.$": "$"
          },
          "End": true
        }
      }
    },
    "ItemReader": {
      "Resource": "arn:aws:states:::s3:listObjectsV2",
      "Parameters": {
        "Bucket": "my-input-bucket",
        "Prefix": "input/"
      }
    },
    "MaxConcurrency": 1000,
    "Next": "Done"
  }
}

ベストプラクティス

設計

  • Standardワークフローを基本とし、高スループットが必要な場合のみExpressを検討する
  • ステートマシンは1つの責務に集中させる(マイクロサービス的に分割)
  • Lambda関数ではなくSDK統合を積極的に使う(コスト削減・シンプル化)
  • 補償トランザクション(Sagaパターン)を使って分散トランザクションを管理する

エラー処理

  • すべてのTaskステートにRetryとCatchを設定する
  • 一時的なエラーには指数バックオフ付きリトライを設定する
  • JitterStrategy: FULLを設定してリトライの集中を避ける
  • States.ALLのCatchを最後に配置する(最後のフォールバック)

運用

  • CloudWatch Logsを有効化してデバッグ情報を記録する
  • X-Rayトレーシングを有効にして全体のレイテンシを把握する
  • 実行履歴の保持期間に注意する(Standard: 90日)
  • ステートマシンのバージョニングにはエイリアスを活用する

コスト

  • 不要な状態遷移を減らす(PassステートやWaitステートの乱用を避ける)
  • 短時間・高頻度の処理にはExpressワークフローを使う
  • SDK統合でLambda関数を削減すると、Lambdaの料金も節約できる

IaC での定義

AWS SAM

Resources:
  OrderStateMachine:
    Type: AWS::Serverless::StateMachine
    Properties:
      DefinitionUri: statemachine/order-processing.asl.json
      DefinitionSubstitutions:
        ValidateOrderFunctionArn: !GetAtt ValidateOrderFunction.Arn
        ProcessPaymentFunctionArn: !GetAtt ProcessPaymentFunction.Arn
      Policies:
        - LambdaInvokePolicy:
            FunctionName: !Ref ValidateOrderFunction
        - LambdaInvokePolicy:
            FunctionName: !Ref ProcessPaymentFunction
      Logging:
        Destinations:
          - CloudWatchLogsLogGroup:
              LogGroupArn: !GetAtt StateMachineLogGroup.Arn
        Level: ALL
        IncludeExecutionData: true

Terraform

resource "aws_sfn_state_machine" "order_processing" {
  name     = "order-processing"
  role_arn = aws_iam_role.sfn_role.arn

  definition = templatefile("${path.module}/statemachine/order-processing.asl.json", {
    validate_order_arn  = aws_lambda_function.validate_order.arn
    process_payment_arn = aws_lambda_function.process_payment.arn
  })

  logging_configuration {
    log_destination        = "${aws_cloudwatch_log_group.sfn.arn}:*"
    include_execution_data = true
    level                  = "ALL"
  }

  tracing_configuration {
    enabled = true
  }
}

参考リンク