CloudFront Origin GroupのMembers入れ替えをStep FunctionsとLambdaで試してみた

CloudFront Origin GroupのMembers入れ替えをStep FunctionsとLambdaで試してみた

CloudFrontのOrigin Groupでフェイルオーバーが継続する状況を想定し、Members順序を入れ替える処理を、Step Functions(JSONata + HTTP Task + SDK統合)とLambdaの2つのアプローチで検証しました。実行契機の自動化は今回の検証範囲に含めていません。結論として現時点では発動させず、CloudFrontのメトリクスを観察して判断する方針です。
2026.07.24

はじめに

前回の記事では、CloudFrontのVPC OriginでOrigin Groupを構成し、応答遅延時に10秒でフェイルオーバーさせる設定を検証しました。

https://dev.classmethod.jp/articles/cloudfront-vpc-origin-failover-timeout-optimization/

フェイルオーバー自体は実現できましたが、Origin Groupのフェイルオーバーはリクエスト単位の迂回であり、Primaryが継続的にダウンしていると毎回接続試行とフェイルオーバーが発生します。そこで今回は、フェイルオーバー発動後にOrigin GroupのMembers順序を入れ替え、セカンダリを新しいプライマリとして扱うための処理を検証しました。

2つのアプローチで検証しています。

  • Step Functions(JSONata + HTTP Task + SDK統合): Lambda不要でヘルスチェックからOrigin Group更新まで完結
  • Lambda: 単一関数にヘルスチェック・判定・更新を集約

結論として、Lambda版でMembers入れ替え処理の動作は確認できましたが、現時点では発動させず、CloudFrontの 5xxErrorRate や追加メトリクスの OriginLatency を観察して必要性を判断する方針です。

テスト環境

両アプローチ共通のテスト環境です。

  • Primary Lambda: Function URL(AuthType: NONE)。Step Functions版では環境変数 RESPONSE_MODE で200/503を切り替え、Lambda版の掲載テンプレートでは DELAY_SECONDS で応答遅延のみを設定可能
  • Secondary Lambda: Function URL(AuthType: NONE)。デフォルトでは遅延なしで200を応答。SecondaryDelaySeconds で応答遅延を設定可能
  • CloudFront Distribution: Origin Group(Primary → Secondary)でフェイルオーバー構成。500/502/503/504でリクエスト単位のフェイルオーバー

Function URLの公開呼び出しには lambda:InvokeFunctionUrl の付与が必要です。2025年10月以降に作成された新しいFunction URLでは、加えてFunction URL経由に限定した lambda:InvokeFunction 権限も必要です。掲載テンプレートでは両方を付与しています。

デプロイは aws cloudformation deploy のみで、SAM Transform不使用の純粋なCloudFormationです。

項目
Step Functions版のリージョナルリソース ap-northeast-1(Lambda, Step Functions, EventBridge Connection)
Lambda版のリージョナルリソース us-east-1(Lambda)
CloudFront Distribution グローバル

アプローチ1: Step Functions(JSONata + HTTP Task + SDK統合)

最初にStep Functionsで実装を試みました。VPCに起因するトラブル(サブネットやENI枯渇など)を避けたかったこと、Lambda実装の複雑化を避けたかったことが動機です。HTTP TaskとSDK統合を使えばLambda不要でヘルスチェックからAPI呼び出しまで完結できます。

処理フロー

Initialize → HealthCheck(Parallel: Primary/Secondary並行)→ CheckLoopComplete → WaitBetweenChecks → HealthCheck(ループバック)→ EvaluateResults →〔条件成立〕GetDistributionConfig → SwapOriginGroupMembers → UpdateDistribution /〔未達〕NoActionNeeded

3回のヘルスチェック完了後、Primary側が3回すべて失敗し、Secondary側が3回すべて成功した場合にMembers入れ替え要求を行います。

ステートマシン定義JSON全文
{
  "Comment": "CloudFront Origin Group - Primary側3回失敗でMembers入れ替え",
  "QueryLanguage": "JSONata",
  "StartAt": "Initialize",
  "States": {
    "Initialize": {
      "Type": "Pass",
      "Comment": "入力値を変数に格納 & カウンター初期化",
      "Assign": {
        "primaryUrl": "{% $states.input.primaryUrl %}",
        "secondaryUrl": "{% $states.input.secondaryUrl %}",
        "connectionArn": "{% $states.input.connectionArn %}",
        "distributionId": "{% $states.input.distributionId %}",
        "counter": 0,
        "primaryFailures": 0,
        "secondaryOK": 0
      },
      "Next": "HealthCheck"
    },
    "HealthCheck": {
      "Type": "Parallel",
      "Comment": "Primary/Secondaryに並行ヘルスチェック",
      "Branches": [
        {
          "StartAt": "CheckPrimary",
          "States": {
            "CheckPrimary": {
              "Type": "Task",
              "Resource": "arn:aws:states:::http:invoke",
              "Arguments": {
                "ApiEndpoint": "{% $primaryUrl %}",
                "Method": "GET",
                "Authentication": {
                  "ConnectionArn": "{% $connectionArn %}"
                }
              },
              "Output": {
                "statusCode": "{% $states.result.StatusCode %}",
                "healthy": true
              },
              "Catch": [
                {
                  "ErrorEquals": ["States.ALL"],
                  "Next": "PrimaryFailed"
                }
              ],
              "End": true
            },
            "PrimaryFailed": {
              "Type": "Pass",
              "Output": {
                "statusCode": 0,
                "healthy": false
              },
              "End": true
            }
          }
        },
        {
          "StartAt": "CheckSecondary",
          "States": {
            "CheckSecondary": {
              "Type": "Task",
              "Resource": "arn:aws:states:::http:invoke",
              "Arguments": {
                "ApiEndpoint": "{% $secondaryUrl %}",
                "Method": "GET",
                "Authentication": {
                  "ConnectionArn": "{% $connectionArn %}"
                }
              },
              "Output": {
                "statusCode": "{% $states.result.StatusCode %}",
                "healthy": true
              },
              "Catch": [
                {
                  "ErrorEquals": ["States.ALL"],
                  "Next": "SecondaryFailed"
                }
              ],
              "End": true
            },
            "SecondaryFailed": {
              "Type": "Pass",
              "Output": {
                "statusCode": 0,
                "healthy": false
              },
              "End": true
            }
          }
        }
      ],
      "Assign": {
        "counter": "{% $counter + 1 %}",
        "primaryFailures": "{% $states.result[0].healthy = false ? $primaryFailures + 1 : $primaryFailures %}",
        "secondaryOK": "{% $states.result[1].healthy = true ? $secondaryOK + 1 : $secondaryOK %}"
      },
      "Next": "CheckLoopComplete"
    },
    "CheckLoopComplete": {
      "Type": "Choice",
      "Comment": "3回チェック完了か判定",
      "Choices": [
        {
          "Condition": "{% $counter >= 3 %}",
          "Next": "EvaluateResults"
        }
      ],
      "Default": "WaitBetweenChecks"
    },
    "WaitBetweenChecks": {
      "Type": "Wait",
      "Seconds": 20,
      "Next": "HealthCheck"
    },
    "EvaluateResults": {
      "Type": "Choice",
      "Comment": "Primary全滅 & Secondary全正常ならMembers入れ替え",
      "Choices": [
        {
          "Condition": "{% $primaryFailures = 3 and $secondaryOK = 3 %}",
          "Next": "GetDistributionConfig"
        }
      ],
      "Default": "NoActionNeeded"
    },
    "NoActionNeeded": {
      "Type": "Succeed",
      "Comment": "入れ替え条件未達、正常終了"
    },
    "GetDistributionConfig": {
      "Type": "Task",
      "Resource": "arn:aws:states:::aws-sdk:cloudfront:getDistributionConfig",
      "Arguments": {
        "Id": "{% $distributionId %}"
      },
      "Assign": {
        "etag": "{% $states.result.ETag %}",
        "distributionConfig": "{% $states.result.DistributionConfig %}"
      },
      "Next": "SwapOriginGroupMembers"
    },
    "SwapOriginGroupMembers": {
      "Type": "Pass",
      "Comment": "Origin GroupのMembers順序を入れ替え(Primary⇔Secondary)",
      "Assign": {
        "updatedConfig": "{% ( $swapped := $distributionConfig.OriginGroups.Items.{ 'Id': Id, 'FailoverCriteria': FailoverCriteria, 'Members': { 'Quantity': Members.Quantity, 'Items': $reverse(Members.Items) } }; $merge([$distributionConfig, { 'OriginGroups': { 'Quantity': $distributionConfig.OriginGroups.Quantity, 'Items': $type($swapped) = 'array' ? $swapped : [$swapped] } }]) ) %}"
      },
      "Next": "UpdateDistribution"
    },
    "UpdateDistribution": {
      "Type": "Task",
      "Resource": "arn:aws:states:::aws-sdk:cloudfront:updateDistribution",
      "Arguments": {
        "Id": "{% $distributionId %}",
        "IfMatch": "{% $etag %}",
        "DistributionConfig": "{% $updatedConfig %}"
      },
      "Output": {
        "status": "UPDATE_ACCEPTED",
        "distributionId": "{% $distributionId %}",
        "etag": "{% $states.result.ETag %}"
      },
      "End": true
    }
  }
}

Initialize: 入力値の変数退避とカウンター初期化

最初のPassステートで、実行時入力(primaryUrl / secondaryUrl / connectionArn / distributionId)を Assign で変数へ格納します。あわせてヘルスチェックのカウンター類を0で初期化します。

"Initialize": {
  "Type": "Pass",
  "Assign": {
    "primaryUrl": "{% $states.input.primaryUrl %}",
    "secondaryUrl": "{% $states.input.secondaryUrl %}",
    "connectionArn": "{% $states.input.connectionArn %}",
    "distributionId": "{% $states.input.distributionId %}",
    "counter": 0,
    "primaryFailures": 0,
    "secondaryOK": 0
  },
  "Next": "HealthCheck"
}

後続のステートでは $states.input ではなく、ここで代入した変数を参照します。Parallelステート通過後に $states.input の構造が変わるためです(詳細は「設計上の注意点」参照)。

HealthCheck: HTTP Taskによる並行ヘルスチェック

Parallelステートで Primary/Secondary を並行チェックします。HTTP Task(arn:aws:states:::http:invoke)でFunction URLにGETリクエストを送り、2xx以外のHTTPレスポンスは States.Http.StatusCode.xxx として失敗します。本実装では States.ALL をCatchしているため、HTTPステータスエラーだけでなく、接続エラーやConnection関連エラーなども一律に healthy: false として扱います。

"CheckPrimary": {
  "Type": "Task",
  "Resource": "arn:aws:states:::http:invoke",
  "Arguments": {
    "ApiEndpoint": "{% $primaryUrl %}",
    "Method": "GET",
    "Authentication": {
      "ConnectionArn": "{% $connectionArn %}"
    }
  },
  "Output": {
    "statusCode": "{% $states.result.StatusCode %}",
    "healthy": true
  },
  "Catch": [
    {
      "ErrorEquals": ["States.ALL"],
      "Next": "PrimaryFailed"
    }
  ],
  "End": true
}

ParallelステートのAssignで各ブランチの結果を集計します。

"Assign": {
  "counter": "{% $counter + 1 %}",
  "primaryFailures": "{% $states.result[0].healthy = false ? $primaryFailures + 1 : $primaryFailures %}",
  "secondaryOK": "{% $states.result[1].healthy = true ? $secondaryOK + 1 : $secondaryOK %}"
}

CheckLoopComplete(Choice)は $counter >= 3 でループ完了を判定し、未達ならWaitBetweenChecks(20秒待機)を経てHealthCheckに戻ります。

SwapOriginGroupMembers: JSONataによるMembers入れ替え

Members入れ替え処理の中核です。JSONataの $reverse$merge を組み合わせてOrigin Group Membersの順序を入れ替えます。

{% (
  $swapped := $distributionConfig.OriginGroups.Items.{
    'Id': Id,
    'FailoverCriteria': FailoverCriteria,
    'Members': {
      'Quantity': Members.Quantity,
      'Items': $reverse(Members.Items)
    }
  };
  $merge([
    $distributionConfig,
    {
      'OriginGroups': {
        'Quantity': $distributionConfig.OriginGroups.Quantity,
        'Items': $type($swapped) = 'array' ? $swapped : [$swapped]
      }
    }
  ])
) %}

UpdateDistribution APIはConfig全体を書き戻す設計のため、$merge で巨大なConfig全体のうちOriginGroupsの部分だけを差し替えています。なお、この式は OriginGroups.Items 全体をマッピングしているため、複数のOrigin Groupが存在するDistributionにそのまま適用すると、すべてのOrigin GroupのMembers順序が反転します。本検証環境ではOrigin Groupが1つだけのため問題ありません。

GetDistributionConfig / UpdateDistribution: SDK統合

"GetDistributionConfig": {
  "Type": "Task",
  "Resource": "arn:aws:states:::aws-sdk:cloudfront:getDistributionConfig",
  "Arguments": {
    "Id": "{% $distributionId %}"
  },
  "Assign": {
    "etag": "{% $states.result.ETag %}",
    "distributionConfig": "{% $states.result.DistributionConfig %}"
  }
}

ETagによる楽観的ロックで、取得後に別の変更が入ると PreconditionFailed で失敗します。

HTTP Task + EventBridge Connection

HTTP TaskにはEventBridge Connectionが必須です。認証不要のエンドポイントに対してもConnectionの作成が必要なため、ダミーのAPI_KEYで作成しています。

HealthCheckConnection:
  Type: AWS::Events::Connection
  Properties:
    Name: origin-failover-healthcheck
    AuthorizationType: API_KEY
    AuthParameters:
      ApiKeyAuthParameters:
        ApiKeyName: x-dummy-key
        ApiKeyValue: dummy-value-for-healthcheck

EventBridge Connectionを作成すると認証情報がSecrets Managerに保存されます。シークレットの保存料(月額$0.40、執筆時点。料金ページ参照)が発生します。

Step Functionsの実行ロールには以下のIAM権限が必要です。

- Effect: Allow
  Action: [ states:InvokeHTTPEndpoint ]
  Resource: "*"
  Condition:
    StringEquals:
      states:HTTPMethod: GET
- Effect: Allow
  Action: [ events:RetrieveConnectionCredentials ]
  Resource: !GetAtt HealthCheckConnection.Arn
- Effect: Allow
  Action: [ secretsmanager:GetSecretValue, secretsmanager:DescribeSecret ]
  Resource: "arn:aws:secretsmanager:*:*:secret:events!connection/*"

設計上の注意点

Parallel内での変数スコープ

JSONataモードではParallelステートの出力は各ブランチの結果配列になります。通過後のステートでは元の入力フィールドを $states.input で参照できません。Initializeで入力値をAssignに格納することで、ループ2回目以降でもParallelの出力に影響されず安定して値を取得できます。

配列のシングルトン問題

JSONataのマッピング(.{...})は結果が1要素の場合にオブジェクトとしてアンラップされます。$type($swapped) = 'array' ? $swapped : [$swapped] で配列を保証し、APIバリデーションエラーを防いでいます。SDK統合でAWSのAPIを呼び出す場面では同様のパターンが必要になることがあります。

Step Functions版の動作確認

Primary正常時(NoActionNeeded)

Primary/Secondaryともに正常応答。EvaluateResultsのDefault分岐によりNoActionNeeded(Succeed)で終了。実行時間は約50秒(Wait 20秒 × 2回 + ヘルスチェック処理時間)。

Primary障害時(Members入れ替え実行)

Primaryを RESPONSE_MODE=unhealthy に設定し、503を返す状態で実行。3回すべてでPrimary側のHTTP Taskが失敗し(States.Http.StatusCode.503States.ALL でCatch)、healthy: false として集計されました。counter=3、primaryFailures=3、secondaryOK=3 で入れ替え条件を満たし、Origin Group Membersの入れ替え要求が成功。

{"status": "UPDATE_ACCEPTED", "distributionId": "EXXXXXXXXXX", "etag": "..."}

UpdateDistribution が正常終了し、Members順序の更新要求が受け付けられました。全エッジロケーションへの反映完了はDistributionのステータスが Deployed になった時点です。本ステートマシンでは反映完了の待機は行っていません。

Primary障害時のStep Functions実行結果(グラフビュー)

CheckPrimaryに「キャッチされたエラー」(⚠️)マークが表示されていますが、503はヘルスチェック失敗としてCatchで処理される想定内の動作です。

Origin Group Members入れ替え確認

タイミング Members順序
実行前 [primary-origin, secondary-origin]
実行後 [secondary-origin, primary-origin]

CloudFront経由でSecondary応答確認

curl https://dXXXXXXXXXXXXX.cloudfront.net/
{"status": "healthy", "origin": "secondary"}

アプローチ2: Lambda

Step Functions版はHTTP TaskやJSONataによるデータ加工など技術的には面白い構成ですが、EventBridge Connectionの必須化(+ Secrets Manager費用)やループ制御の冗長さが気になりました。

当初LambdaはVPC関連のトラブル(サブネット設定やENI枯渇など)を避けたくて選択肢から外していましたが、このLambdaはCloudFront APIとオリジンへのHTTPS通信しか行わないためVPC配置が不要です。VPCの懸念がなくなったため、シンプルにLambda単体で書き直しました。

処理フロー

  1. GetDistributionConfig で Origin Group の現在のMembers順序を取得
  2. セカンダリ(2番目のメンバー)のドメインに HTTPS ヘルスチェック(10秒間隔 × 5回)
  3. 平均レイテンシ ≤ 閾値 かつ エラー率 ≤ 閾値 → 「セカンダリへの直接HTTPS GETが正常」と判定
  4. しきい値を満たした場合は、Members順を入れ替えて UpdateDistribution を呼び出す
  5. しきい値を満たさない場合は何もしない(直接ヘルスチェックが失敗しているセカンダリへの切り替えを抑止)

Step Functions版との設計上の違いとして、Lambda版はセカンダリのみをヘルスチェックします。今回のLambda版では「プライマリに障害が発生している」という呼び出し側の前提を置き、切り替え先のドメインに対する直接HTTPS GETのみを判定対象としました。必須でない要素を削り、Lambda実装の複雑化を避けています。

コード要点

ヘルスチェック部分:

def health_check_once(domain, path, timeout=10):
    url = f"https://{domain}{path}"
    ctx = ssl.create_default_context()
    start = time.time()
    try:
        req = urllib.request.Request(url, method="GET")
        req.add_header("User-Agent", "CloudFront-Failover-HealthCheck/1.0")
        with urllib.request.urlopen(req, timeout=timeout, context=ctx) as resp:
            latency = time.time() - start
            return {"success": 200 <= resp.status < 400, "status_code": resp.status, "latency": round(latency, 3)}
    except urllib.error.HTTPError as e:
        latency = time.time() - start
        return {"success": False, "status_code": e.code, "latency": round(latency, 3)}
    except Exception as e:
        latency = time.time() - start
        return {"success": False, "status_code": 0, "latency": round(latency, 3), "error": str(e)}

判定部分:

def evaluate_health(results, max_latency, max_error_rate):
    total = len(results)
    errors = sum(1 for r in results if not r["success"])
    error_rate = errors / total
    avg_latency = sum(r["latency"] for r in results) / total
    healthy = avg_latency <= max_latency and error_rate <= max_error_rate
    return healthy, f"avg_latency={avg_latency:.3f}s, error_rate={error_rate:.0%}"

Origin Group Members入れ替え部分:

def swap_origin_group_members(distribution_config, origin_group_id):
    origin_groups = distribution_config.get("OriginGroups", {})
    for group in origin_groups.get("Items", []):
        if group["Id"] == origin_group_id:
            members = group["Members"]["Items"]
            group["Members"]["Items"] = [members[1], members[0]]
            return members[0]["OriginId"], members[1]["OriginId"]
    raise ValueError(f"Origin Group '{origin_group_id}' not found")

Step Functions版の $reverse + $merge と同じことを、Pythonのリスト操作で直接的に実現しています。Distribution Config全体を update_distribution に渡す点は同じです。

Lambda関数コード全文
import json
import os
import time
import urllib.request
import urllib.error
import ssl

import boto3

def get_config():
    return {
        "distribution_id": os.environ["DISTRIBUTION_ID"],
        "origin_group_id": os.environ.get("ORIGIN_GROUP_ID", "failover-group"),
        "check_count": int(os.environ.get("CHECK_COUNT", "5")),
        "check_interval": int(os.environ.get("CHECK_INTERVAL", "10")),
        "max_latency": float(os.environ.get("MAX_LATENCY", "1.0")),
        "max_error_rate": float(os.environ.get("MAX_ERROR_RATE", "0.0")),
        "health_check_path": os.environ.get("HEALTH_CHECK_PATH", "/"),
    }

def get_distribution_config(client, distribution_id):
    resp = client.get_distribution_config(Id=distribution_id)
    return resp["DistributionConfig"], resp["ETag"]

def find_origin_group(distribution_config, origin_group_id):
    origin_groups = distribution_config.get("OriginGroups", {})
    for group in origin_groups.get("Items", []):
        if group["Id"] == origin_group_id:
            return group
    raise ValueError(f"Origin Group '{origin_group_id}' not found")

def get_origin_domain(distribution_config, origin_id):
    for origin in distribution_config["Origins"]["Items"]:
        if origin["Id"] == origin_id:
            return origin["DomainName"]
    raise ValueError(f"Origin '{origin_id}' not found")

def health_check_once(domain, path, timeout=10):
    url = f"https://{domain}{path}"
    ctx = ssl.create_default_context()
    start = time.time()
    try:
        req = urllib.request.Request(url, method="GET")
        req.add_header("User-Agent", "CloudFront-Failover-HealthCheck/1.0")
        with urllib.request.urlopen(req, timeout=timeout, context=ctx) as resp:
            latency = time.time() - start
            return {"success": 200 <= resp.status < 400, "status_code": resp.status, "latency": round(latency, 3)}
    except urllib.error.HTTPError as e:
        latency = time.time() - start
        return {"success": False, "status_code": e.code, "latency": round(latency, 3)}
    except Exception as e:
        latency = time.time() - start
        return {"success": False, "status_code": 0, "latency": round(latency, 3), "error": str(e)}

def run_health_checks(domain, path, count, interval):
    results = []
    for i in range(count):
        if i > 0:
            time.sleep(interval)
        result = health_check_once(domain, path)
        results.append(result)
        print(f"  Check {i+1}/{count}: status={result.get('status_code')} latency={result['latency']}s success={result['success']}")
    return results

def evaluate_health(results, max_latency, max_error_rate):
    if not results:
        return False, "No results"
    total = len(results)
    errors = sum(1 for r in results if not r["success"])
    error_rate = errors / total
    avg_latency = sum(r["latency"] for r in results) / total
    summary = f"avg_latency={avg_latency:.3f}s (max={max_latency}s), error_rate={error_rate:.0%} (max={max_error_rate:.0%})"
    healthy = avg_latency <= max_latency and error_rate <= max_error_rate
    return healthy, summary

def swap_origin_group_members(distribution_config, origin_group_id):
    origin_groups = distribution_config.get("OriginGroups", {})
    for group in origin_groups.get("Items", []):
        if group["Id"] == origin_group_id:
            members = group["Members"]["Items"]
            if len(members) != 2:
                raise ValueError(f"Origin Group '{origin_group_id}' has {len(members)} members, expected 2")
            group["Members"]["Items"] = [members[1], members[0]]
            return members[0]["OriginId"], members[1]["OriginId"]
    raise ValueError(f"Origin Group '{origin_group_id}' not found")

def handler(event, context):
    config = get_config()
    print(f"Config: distribution_id={config['distribution_id']}, origin_group_id={config['origin_group_id']}")

    cf_client = boto3.client("cloudfront")
    distribution_config, etag = get_distribution_config(cf_client, config["distribution_id"])

    origin_group = find_origin_group(distribution_config, config["origin_group_id"])
    members = origin_group["Members"]["Items"]
    primary_id = members[0]["OriginId"]
    secondary_id = members[1]["OriginId"]

    primary_domain = get_origin_domain(distribution_config, primary_id)
    secondary_domain = get_origin_domain(distribution_config, secondary_id)

    print(f"Current order:")
    print(f"  Primary:   {primary_id} ({primary_domain})")
    print(f"  Secondary: {secondary_id} ({secondary_domain})")

    print(f"\nHealth checking secondary: {secondary_domain}")
    results = run_health_checks(
        domain=secondary_domain,
        path=config["health_check_path"],
        count=config["check_count"],
        interval=config["check_interval"],
    )

    healthy, summary = evaluate_health(results, config["max_latency"], config["max_error_rate"])
    print(f"\nSecondary direct check: {'PASSED' if healthy else 'FAILED'} - {summary}")

    if not healthy:
        msg = f"Secondary direct check failed. No member swap requested. {summary}"
        print(msg)
        return {
            "action": "none",
            "reason": "secondary_direct_check_failed",
            "summary": summary,
            "primary": primary_id,
            "secondary": secondary_id,
        }

    print(f"\nSwapping origin group members: {primary_id} <-> {secondary_id}")
    old_primary, old_secondary = swap_origin_group_members(distribution_config, config["origin_group_id"])

    new_etag = cf_client.update_distribution(
        Id=config["distribution_id"],
        IfMatch=etag,
        DistributionConfig=distribution_config,
    )["ETag"]

    print(f"Distribution update accepted. New ETag: {new_etag}")
    print(f"Requested order: Primary={old_secondary}, Secondary={old_primary}")

    return {
        "action": "swapped",
        "old_primary": old_primary,
        "old_secondary": old_secondary,
        "new_primary": old_secondary,
        "new_secondary": old_primary,
        "health_summary": summary,
    }
CloudFormationテンプレート全文(テスト環境一式)
AWSTemplateFormatVersion: '2010-09-09'
Description: >-
  Test environment for CloudFront Origin Group Failover Lambda.
  Creates a CloudFront distribution with origin group (failover-group)
  using Lambda Function URLs as primary/secondary origins.

Parameters:
  ProjectName:
    Type: String
    Default: cf-failover-test

  PrimaryDelaySeconds:
    Type: Number
    Default: 0
    Description: Primary Lambda response delay in seconds

  SecondaryDelaySeconds:
    Type: Number
    Default: 0
    Description: Secondary Lambda response delay in seconds (keep 0 for healthy secondary)

  HealthCheckCount:
    Type: Number
    Default: 5
    Description: Number of health checks the failover Lambda performs

  HealthCheckInterval:
    Type: Number
    Default: 10
    Description: Interval between health checks in seconds

Resources:
  OriginLambdaRole:
    Type: AWS::IAM::Role
    Properties:
      RoleName: !Sub ${ProjectName}-origin-role
      AssumeRolePolicyDocument:
        Version: '2012-10-17'
        Statement:
          - Effect: Allow
            Principal:
              Service: lambda.amazonaws.com
            Action: sts:AssumeRole
      ManagedPolicyArns:
        - arn:aws:iam::aws:policy/service-role/AWSLambdaBasicExecutionRole

  PrimaryFunction:
    Type: AWS::Lambda::Function
    Properties:
      FunctionName: !Sub ${ProjectName}-primary
      Runtime: python3.12
      Handler: index.handler
      Role: !GetAtt OriginLambdaRole.Arn
      Timeout: 30
      MemorySize: 128
      Environment:
        Variables:
          DELAY_SECONDS: !Ref PrimaryDelaySeconds
      Code:
        ZipFile: |
          import json, os, time
          def handler(event, context):
              delay = int(os.environ.get("DELAY_SECONDS", "0"))
              if delay > 0:
                  time.sleep(delay)
              return {
                  "statusCode": 200,
                  "headers": {"Content-Type": "application/json"},
                  "body": json.dumps({"origin": "primary", "delay": delay})
              }

  PrimaryFunctionUrl:
    Type: AWS::Lambda::Url
    Properties:
      AuthType: NONE
      TargetFunctionArn: !GetAtt PrimaryFunction.Arn

  PrimaryFunctionUrlPermission:
    Type: AWS::Lambda::Permission
    Properties:
      FunctionName: !GetAtt PrimaryFunction.Arn
      Action: lambda:InvokeFunctionUrl
      Principal: '*'
      FunctionUrlAuthType: NONE

  PrimaryFunctionInvokePermission:
    Type: AWS::Lambda::Permission
    Properties:
      FunctionName: !GetAtt PrimaryFunction.Arn
      Action: lambda:InvokeFunction
      Principal: '*'
      InvokedViaFunctionUrl: true

  SecondaryFunction:
    Type: AWS::Lambda::Function
    Properties:
      FunctionName: !Sub ${ProjectName}-secondary
      Runtime: python3.12
      Handler: index.handler
      Role: !GetAtt OriginLambdaRole.Arn
      Timeout: 30
      MemorySize: 128
      Environment:
        Variables:
          DELAY_SECONDS: !Ref SecondaryDelaySeconds
      Code:
        ZipFile: |
          import json, os, time
          def handler(event, context):
              delay = int(os.environ.get("DELAY_SECONDS", "0"))
              if delay > 0:
                  time.sleep(delay)
              return {
                  "statusCode": 200,
                  "headers": {"Content-Type": "application/json"},
                  "body": json.dumps({"origin": "secondary", "delay": delay})
              }

  SecondaryFunctionUrl:
    Type: AWS::Lambda::Url
    Properties:
      AuthType: NONE
      TargetFunctionArn: !GetAtt SecondaryFunction.Arn

  SecondaryFunctionUrlPermission:
    Type: AWS::Lambda::Permission
    Properties:
      FunctionName: !GetAtt SecondaryFunction.Arn
      Action: lambda:InvokeFunctionUrl
      Principal: '*'
      FunctionUrlAuthType: NONE

  SecondaryFunctionInvokePermission:
    Type: AWS::Lambda::Permission
    Properties:
      FunctionName: !GetAtt SecondaryFunction.Arn
      Action: lambda:InvokeFunction
      Principal: '*'
      InvokedViaFunctionUrl: true

  Distribution:
    Type: AWS::CloudFront::Distribution
    Properties:
      DistributionConfig:
        Enabled: true
        Comment: !Sub ${ProjectName} - Origin Group Failover Test
        DefaultCacheBehavior:
          TargetOriginId: failover-group
          ViewerProtocolPolicy: redirect-to-https
          CachePolicyId: 4135ea2d-6df8-44a3-9df3-4b5a84be39ad
          OriginRequestPolicyId: b689b0a8-53d0-40ab-baf2-68738e2966ac
          AllowedMethods: [GET, HEAD]
          Compress: true
        Origins:
          - Id: primary-origin
            DomainName: !Select [2, !Split ['/', !GetAtt PrimaryFunctionUrl.FunctionUrl]]
            CustomOriginConfig:
              HTTPSPort: 443
              OriginProtocolPolicy: https-only
              OriginSSLProtocols: [TLSv1.2]
          - Id: secondary-origin
            DomainName: !Select [2, !Split ['/', !GetAtt SecondaryFunctionUrl.FunctionUrl]]
            CustomOriginConfig:
              HTTPSPort: 443
              OriginProtocolPolicy: https-only
              OriginSSLProtocols: [TLSv1.2]
        OriginGroups:
          Quantity: 1
          Items:
            - Id: failover-group
              FailoverCriteria:
                StatusCodes:
                  Quantity: 4
                  Items: [500, 502, 503, 504]
              Members:
                Quantity: 2
                Items:
                  - OriginId: primary-origin
                  - OriginId: secondary-origin
        HttpVersion: http2and3
        PriceClass: PriceClass_100

  FailoverLambdaRole:
    Type: AWS::IAM::Role
    Properties:
      RoleName: !Sub ${ProjectName}-failover-role
      AssumeRolePolicyDocument:
        Version: '2012-10-17'
        Statement:
          - Effect: Allow
            Principal:
              Service: lambda.amazonaws.com
            Action: sts:AssumeRole
      ManagedPolicyArns:
        - arn:aws:iam::aws:policy/service-role/AWSLambdaBasicExecutionRole
      Policies:
        - PolicyName: CloudFrontAccess
          PolicyDocument:
            Version: '2012-10-17'
            Statement:
              - Effect: Allow
                Action:
                  - cloudfront:GetDistribution
                  - cloudfront:GetDistributionConfig
                  - cloudfront:UpdateDistribution
                Resource: !Sub arn:aws:cloudfront::${AWS::AccountId}:distribution/${Distribution}

  FailoverFunction:
    Type: AWS::Lambda::Function
    Properties:
      FunctionName: !Sub ${ProjectName}-failover
      Runtime: python3.12
      Handler: index.handler
      Role: !GetAtt FailoverLambdaRole.Arn
      Timeout: 120
      MemorySize: 128
      Environment:
        Variables:
          DISTRIBUTION_ID: !Ref Distribution
          ORIGIN_GROUP_ID: failover-group
          CHECK_COUNT: !Ref HealthCheckCount
          CHECK_INTERVAL: !Ref HealthCheckInterval
          MAX_LATENCY: '1.0'
          MAX_ERROR_RATE: '0.0'
          HEALTH_CHECK_PATH: /
      Code:
        ZipFile: |
          import json, os, time, urllib.request, urllib.error, ssl
          import boto3

          def get_config():
              return {
                  "distribution_id": os.environ["DISTRIBUTION_ID"],
                  "origin_group_id": os.environ.get("ORIGIN_GROUP_ID", "failover-group"),
                  "check_count": int(os.environ.get("CHECK_COUNT", "5")),
                  "check_interval": int(os.environ.get("CHECK_INTERVAL", "10")),
                  "max_latency": float(os.environ.get("MAX_LATENCY", "1.0")),
                  "max_error_rate": float(os.environ.get("MAX_ERROR_RATE", "0.0")),
                  "health_check_path": os.environ.get("HEALTH_CHECK_PATH", "/"),
              }

          def get_distribution_config(client, distribution_id):
              resp = client.get_distribution_config(Id=distribution_id)
              return resp["DistributionConfig"], resp["ETag"]

          def find_origin_group(distribution_config, origin_group_id):
              for group in distribution_config.get("OriginGroups", {}).get("Items", []):
                  if group["Id"] == origin_group_id:
                      return group
              raise ValueError(f"Origin Group '{origin_group_id}' not found")

          def get_origin_domain(distribution_config, origin_id):
              for origin in distribution_config["Origins"]["Items"]:
                  if origin["Id"] == origin_id:
                      return origin["DomainName"]
              raise ValueError(f"Origin '{origin_id}' not found")

          def health_check_once(domain, path, timeout=10):
              url = f"https://{domain}{path}"
              ctx = ssl.create_default_context()
              start = time.time()
              try:
                  req = urllib.request.Request(url, method="GET")
                  req.add_header("User-Agent", "CloudFront-Failover-HealthCheck/1.0")
                  with urllib.request.urlopen(req, timeout=timeout, context=ctx) as resp:
                      latency = time.time() - start
                      return {"success": 200 <= resp.status < 400, "status_code": resp.status, "latency": round(latency, 3)}
              except urllib.error.HTTPError as e:
                  latency = time.time() - start
                  return {"success": False, "status_code": e.code, "latency": round(latency, 3)}
              except Exception as e:
                  latency = time.time() - start
                  return {"success": False, "status_code": 0, "latency": round(latency, 3), "error": str(e)}

          def run_health_checks(domain, path, count, interval):
              results = []
              for i in range(count):
                  if i > 0:
                      time.sleep(interval)
                  result = health_check_once(domain, path)
                  results.append(result)
                  print(f"  Check {i+1}/{count}: status={result.get('status_code')} latency={result['latency']}s success={result['success']}")
              return results

          def evaluate_health(results, max_latency, max_error_rate):
              if not results:
                  return False, "No results"
              total = len(results)
              errors = sum(1 for r in results if not r["success"])
              error_rate = errors / total
              avg_latency = sum(r["latency"] for r in results) / total
              summary = f"avg_latency={avg_latency:.3f}s (max={max_latency}s), error_rate={error_rate:.0%} (max={max_error_rate:.0%})"
              healthy = avg_latency <= max_latency and error_rate <= max_error_rate
              return healthy, summary

          def swap_origin_group_members(distribution_config, origin_group_id):
              for group in distribution_config.get("OriginGroups", {}).get("Items", []):
                  if group["Id"] == origin_group_id:
                      members = group["Members"]["Items"]
                      if len(members) != 2:
                          raise ValueError(f"Expected 2 members, got {len(members)}")
                      group["Members"]["Items"] = [members[1], members[0]]
                      return members[0]["OriginId"], members[1]["OriginId"]
              raise ValueError(f"Origin Group '{origin_group_id}' not found")

          def handler(event, context):
              config = get_config()
              print(f"Config: distribution_id={config['distribution_id']}, origin_group_id={config['origin_group_id']}")
              cf_client = boto3.client("cloudfront")
              distribution_config, etag = get_distribution_config(cf_client, config["distribution_id"])
              origin_group = find_origin_group(distribution_config, config["origin_group_id"])
              members = origin_group["Members"]["Items"]
              primary_id = members[0]["OriginId"]
              secondary_id = members[1]["OriginId"]
              primary_domain = get_origin_domain(distribution_config, primary_id)
              secondary_domain = get_origin_domain(distribution_config, secondary_id)
              print(f"Current order:\n  Primary:   {primary_id} ({primary_domain})\n  Secondary: {secondary_id} ({secondary_domain})")
              print(f"\nHealth checking secondary: {secondary_domain}")
              results = run_health_checks(secondary_domain, config["health_check_path"], config["check_count"], config["check_interval"])
              healthy, summary = evaluate_health(results, config["max_latency"], config["max_error_rate"])
              print(f"\nSecondary direct check: {'PASSED' if healthy else 'FAILED'} - {summary}")
              if not healthy:
                  print(f"Secondary direct check failed. No member swap requested. {summary}")
                  return {"action": "none", "reason": "secondary_direct_check_failed", "summary": summary}
              print(f"\nSwapping origin group members: {primary_id} <-> {secondary_id}")
              old_primary, old_secondary = swap_origin_group_members(distribution_config, config["origin_group_id"])
              new_etag = cf_client.update_distribution(Id=config["distribution_id"], IfMatch=etag, DistributionConfig=distribution_config)["ETag"]
              print(f"Distribution update accepted. New ETag: {new_etag}")
              print(f"Requested order: Primary={old_secondary}, Secondary={old_primary}")
              return {"action": "swapped", "old_primary": old_primary, "new_primary": old_secondary, "health_summary": summary}

Lambda版の動作確認

CloudFormationでテスト環境一式(Primary/Secondary Lambda + CloudFront + Failover Lambda)をデプロイし、CLIから実行しました。

./invoke.sh cf-failover-test-failover --region us-east-1

実行ログ:

Config: distribution_id=E347LRBJZPL0WN, origin_group_id=failover-group
Current order:
  Primary:   primary-origin (wa72nucfq54qzrfydeupwye3ke0itoqo.lambda-url.us-east-1.on.aws)
  Secondary: secondary-origin (iupokdzzmtr72rd3dryjhqaswy0sfden.lambda-url.us-east-1.on.aws)

Health checking secondary: iupokdzzmtr72rd3dryjhqaswy0sfden.lambda-url.us-east-1.on.aws
  Check 1/5: status=200 latency=0.369s success=True
  Check 2/5: status=200 latency=0.069s success=True
  Check 3/5: status=200 latency=0.079s success=True
  Check 4/5: status=200 latency=0.087s success=True
  Check 5/5: status=200 latency=0.082s success=True

Secondary direct check: PASSED - avg_latency=0.137s (max=1.0s), error_rate=0% (max=0%)

Swapping origin group members: primary-origin <-> secondary-origin
Distribution update accepted. New ETag: ETVPDKIKX0DER
Requested order: Primary=secondary-origin, Secondary=primary-origin

ヘルスチェック5回すべて成功(平均レイテンシ0.137秒)、Origin Group Membersの入れ替え要求が成功しました。実行時間は48.5秒(ヘルスチェック10秒間隔×4 + CloudFront API呼び出し)。掲載のデフォルト設定では、各チェックが約10秒でタイムアウトすると仮定した単純計算で、待機時間と合わせてヘルスチェックに約90秒かかります。後続のCloudFront API呼び出し時間も考慮し、Lambdaのタイムアウトは120秒に設定しています。

入れ替え後のOrigin Group状態:

タイミング Members順序
実行前 [primary-origin, secondary-origin]
実行後 [secondary-origin, primary-origin]

2つのアプローチの比較

観点 Step Functions Lambda
ヘルスチェック対象 Primary + Secondary 両方 Secondary のみ
判定ロジック Primary全滅 & Secondary全正常 Secondary の平均レイテンシ + エラー率
外部依存 EventBridge Connection(Secrets Manager費用) なし
コード管理 JSONベースのステートマシン定義 Python単一ファイル
実行時間 約50秒(20秒wait × 2 + チェック) 約48秒(10秒interval × 4 + チェック)
可読性 JSONata式が複雑(特にswap部分) Pythonで直感的
ランタイム管理 不要 Python 3.12 のEOL管理が必要
拡張性 ステート追加で分岐・通知が容易 コード変更が必要

Step Functions版はLambda不要で構成できる点が魅力ですが、EventBridge Connectionの必須化やJSONataの配列シングルトン問題への対処など、実装上のハードルがあります。Lambda版はPythonで素直に書けるため保守性が高く、今回の用途では十分です。

まとめ

CloudFrontのOrigin Group Members入れ替えを、Step FunctionsとLambdaの2つのアプローチで検証しました。

Lambda版で、セカンダリのヘルスチェックからMembers入れ替え要求までが動作することを確認できました。ただし、実行契機(CloudWatch Alarm連携など)を含む完全な自動化は未実装であり、現時点ではこの処理を自動起動する予定もありません。

  • stale-if-error / stale-while-revalidate によるキャッシュ保護と、OriginReadTimeout 短縮で同期待ち時間を最小化する対策を先行して入れている
  • Origin Groupのリクエスト単位フェイルオーバー(500/502/503/504)で、個々のリクエストレベルではセカンダリへの迂回が機能している

https://dev.classmethod.jp/articles/cloudfront-stale-if-error-origin-timeout-behavior/

https://dev.classmethod.jp/articles/cloudfront-vpc-origin-failover-timeout-optimization/

これらの対策でカバーできない障害が再発した場合に、Origin Group Membersの自動入れ替えを再検討します。障害時にCloudFrontメトリクスの 5xxErrorRate や追加メトリクスの OriginLatency へどのような影響が現れるかを観察し、自動化の必要性やしきい値の実現性を判断します。ただし、5xxErrorRate はビューワーへ返したレスポンスの5xx割合であり、Origin Groupのフェイルオーバーでセカンダリが200を返した場合にはプライマリの5xxが直接現れない点に注意が必要です。実装は今回のLambda版をベースに、起動トリガーやリトライ制御はStep Functionsで組み合わせるハイブリッド構成を想定しています。

参考リンク

https://docs.aws.amazon.com/step-functions/latest/dg/connect-third-party-apis.html

https://docs.aws.amazon.com/step-functions/latest/dg/supported-services-awssdk.html

https://docs.aws.amazon.com/AmazonCloudFront/latest/APIReference/API_GetDistributionConfig.html

https://docs.aws.amazon.com/AmazonCloudFront/latest/APIReference/API_UpdateDistribution.html

この記事をシェアする

AWSのお困り事はクラスメソッドへ

関連記事