CloudFront Origin GroupのMembers入れ替えをStep FunctionsとLambdaで試してみた
はじめに
前回の記事では、CloudFrontのVPC OriginでOrigin Groupを構成し、応答遅延時に10秒でフェイルオーバーさせる設定を検証しました。
フェイルオーバー自体は実現できましたが、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.503 → States.ALL でCatch)、healthy: false として集計されました。counter=3、primaryFailures=3、secondaryOK=3 で入れ替え条件を満たし、Origin Group Membersの入れ替え要求が成功。
{"status": "UPDATE_ACCEPTED", "distributionId": "EXXXXXXXXXX", "etag": "..."}
UpdateDistribution が正常終了し、Members順序の更新要求が受け付けられました。全エッジロケーションへの反映完了はDistributionのステータスが Deployed になった時点です。本ステートマシンでは反映完了の待機は行っていません。

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単体で書き直しました。
処理フロー
GetDistributionConfigで Origin Group の現在のMembers順序を取得- セカンダリ(2番目のメンバー)のドメインに HTTPS ヘルスチェック(10秒間隔 × 5回)
- 平均レイテンシ ≤ 閾値 かつ エラー率 ≤ 閾値 → 「セカンダリへの直接HTTPS GETが正常」と判定
- しきい値を満たした場合は、Members順を入れ替えて
UpdateDistributionを呼び出す - しきい値を満たさない場合は何もしない(直接ヘルスチェックが失敗しているセカンダリへの切り替えを抑止)
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)で、個々のリクエストレベルではセカンダリへの迂回が機能している
これらの対策でカバーできない障害が再発した場合に、Origin Group Membersの自動入れ替えを再検討します。障害時にCloudFrontメトリクスの 5xxErrorRate や追加メトリクスの OriginLatency へどのような影響が現れるかを観察し、自動化の必要性やしきい値の実現性を判断します。ただし、5xxErrorRate はビューワーへ返したレスポンスの5xx割合であり、Origin Groupのフェイルオーバーでセカンダリが200を返した場合にはプライマリの5xxが直接現れない点に注意が必要です。実装は今回のLambda版をベースに、起動トリガーやリトライ制御はStep Functionsで組み合わせるハイブリッド構成を想定しています。
参考リンク








