📁 完全なコードはGitHubで公開:GitHub: pipeiac02
シリーズ一覧(全12回)
Phase 1: 基盤構築
Phase 2: ワークフロー
Phase 3: セキュリティ・運用
Phase 4: 開発効率化
1. はじめに
1-1. 今回のゴール
Step Functionsにエラーハンドリングを追加して、障害に強いパイプラインにします。
| ゴール | 内容 |
|---|---|
| 理解 | エラー処理のパターンを理解する |
| 実装 | Retry/Catch を追加する |
| 通知 | エラー時にSNSで通知する |
1-2. なぜエラー対策が必要か
現状のパイプラインは、エラーが発生すると そのまま失敗 します。
これを リトライ+通知 する仕組みに改善します。
1-3. この記事で作るもの
2. 比喩で理解する
2-1. エラーハンドリングを「料理失敗時の対応」で考える
料理が失敗しても、すぐ諦めるのではなく、再挑戦や報告の仕組みが必要です。
2-2. 比喩の図解(レストラン)
2-3. 対応マニュアルの例
| 状況 | 対応 | 回数 |
|---|---|---|
| 焦げた | 作り直す | 最大3回 |
| 材料切れ | 店長に報告 | 即座に |
| 3回失敗 | 店長に報告して中止 | - |
2-4. 技術の図解(AWS)
2-5. 対応関係
| レストラン | AWS | 役割 |
|---|---|---|
| 作り直し | Retry | 自動リトライ |
| 諦める | Catch | エラーをキャッチ |
| 店長に報告 | SNS通知 | 担当者に通知 |
| 対応マニュアル | エラーハンドリング設定 | 障害対応ルール |
3. Step Functionsのエラー処理
3-1. 2つのパターン
| パターン | 用途 | 料理で例えると |
|---|---|---|
| Retry | 一時的なエラーで再試行 | 焦げたから作り直す |
| Catch | リトライ後も失敗で別処理 | 3回失敗したら店長に報告 |
3-2. Retry の設定項目
| 設定 | 説明 | 例 |
|---|---|---|
| ErrorEquals | リトライ対象のエラー |
States.ALL で全エラー |
| IntervalSeconds | 最初の待機秒数 |
2 秒 |
| MaxAttempts | 最大リトライ回数 |
3 回 |
| BackoffRate | 待機時間の増加率 |
2.0 で倍増 |
3-3. リトライの動作例
3-4. Catch の設定項目
| 設定 | 説明 | 例 |
|---|---|---|
| ErrorEquals | キャッチ対象のエラー |
States.ALL で全エラー |
| Next | 遷移先のステート | NotifyError |
| ResultPath | エラー情報の格納先 | $.error |
3-5. エラーの種類
| エラー | 説明 | 対処 |
|---|---|---|
| States.ALL | 全てのエラー | 汎用的なキャッチ |
| States.Timeout | タイムアウト | リトライ |
| States.TaskFailed | タスク失敗 | リトライ |
| Lambda.ServiceException | Lambda側エラー | リトライ |
| States.Permissions | 権限エラー | 通知(リトライ不要) |
4. 実装
4-1. ファイル構成
今回作成・修正するファイル:
pipeiac02/
├── test-data/
│ └── error-test-ok.json # ★テスト用サンプルデータ
└── tf/
├── stepfunctions.tf # ★Retry/Catch追加
├── sns.tf # ★今回作成
└── iam.tf # ★SNS発行権限追加
4-2. SNSトピック作成(sns.tf)
エラー通知用のSNSトピックを作成します。
フルコード:
# sns.tf
# ================================
# SNS(通知)
# ================================
# SNS Topic(アラート通知用)
# パイプライン失敗時やLambdaエラー時の通知先
resource "aws_sns_topic" "alert" {
name = "${var.project}-alert"
tags = merge(local.common_tags, {
Name = "${var.project}-alert"
})
}
# SNS Subscription(メール通知)
# 指定したメールアドレスに通知を送信
resource "aws_sns_topic_subscription" "email" {
topic_arn = aws_sns_topic.alert.arn
protocol = "email"
endpoint = var.alert_email
}
4-3. IAM権限追加(iam.tf に追加)
Step FunctionsがSNSに発行できる権限を追加します。
フルコード:
# iam.tf に追加
# ================================
# Step Functions → SNS 発行権限
# ================================
resource "aws_iam_role_policy" "sfn_sns" {
name = "${var.project}-sfn-sns-policy"
role = aws_iam_role.sfn.id
policy = jsonencode({
Version = "2012-10-17"
Statement = [
{
Effect = "Allow"
Action = [
"sns:Publish"
]
Resource = aws_sns_topic.alert.arn
}
]
})
}
4-4. Step Functions改善(stepfunctions.tf)
Retry/Catch/通知ステートを追加します。
コードの骨格:
# stepfunctions.tf
States = {
ETL = {
# ... 既存の設定 ...
Retry = [...] # ★追加
Catch = [...] # ★追加
}
StartCrawler = {
# ... 既存の設定 ...
Retry = [...] # ★追加
Catch = [...] # ★追加
}
NotifyError = { # ★新規追加
# SNS通知
}
}
フルコード:
# stepfunctions.tf
# ================================
# Step Functions ステートマシン(エラーハンドリング版)
# ================================
resource "aws_sfn_state_machine" "pipeline" {
name = "${var.project}-pipeline"
role_arn = aws_iam_role.sfn.arn
definition = jsonencode({
Comment = "データパイプライン: ETL → Crawler(エラーハンドリング付き)"
StartAt = "ETL"
States = {
# ================================
# ETL処理(Lambda実行)
# ================================
ETL = {
Type = "Task"
Resource = "arn:aws:states:::lambda:invoke"
Parameters = {
FunctionName = aws_lambda_function.etl.arn
Payload = {
"input_key.$" = "$.input_key"
}
}
ResultPath = "$.etl_result"
# ★リトライ設定
Retry = [
{
ErrorEquals = [
"Lambda.ServiceException",
"Lambda.AWSLambdaException",
"Lambda.SdkClientException",
"States.Timeout"
]
IntervalSeconds = 2
MaxAttempts = 3
BackoffRate = 2.0
}
]
# ★エラーキャッチ
Catch = [
{
ErrorEquals = ["States.ALL"]
ResultPath = "$.error"
Next = "NotifyError"
}
]
Next = "StartCrawler"
}
# ================================
# Crawler起動
# ================================
StartCrawler = {
Type = "Task"
Resource = "arn:aws:states:::aws-sdk:glue:startCrawler"
Parameters = {
Name = aws_glue_crawler.main.name
}
ResultPath = "$.crawler_result"
# ★リトライ設定
Retry = [
{
ErrorEquals = [
"Glue.CrawlerRunningException",
"States.Timeout"
]
IntervalSeconds = 30
MaxAttempts = 3
BackoffRate = 1.5
}
]
# ★エラーキャッチ
Catch = [
{
ErrorEquals = ["States.ALL"]
ResultPath = "$.error"
Next = "NotifyError"
}
]
End = true
}
# ================================
# ★エラー通知(新規追加)
# ================================
NotifyError = {
Type = "Task"
Resource = "arn:aws:states:::sns:publish"
Parameters = {
TopicArn = aws_sns_topic.alert.arn
Subject = "【エラー】データパイプライン失敗"
Message = {
"FunctionName.$" = "$$.Execution.Name"
"Error.$" = "$.error"
"Timestamp.$" = "$$.State.EnteredTime"
}
}
End = true
}
}
})
tags = local.common_tags
}
4-5. コード解説
4-5-1. ETLステートのRetry
Retry = [
{
ErrorEquals = ["Lambda.ServiceException", ...]
IntervalSeconds = 2 # 最初は2秒待機
MaxAttempts = 3 # 最大3回リトライ
BackoffRate = 2.0 # 待機時間を2倍ずつ増加
}
]
| 回数 | 待機時間 |
|---|---|
| 1回目 | 2秒 |
| 2回目 | 4秒(2×2.0) |
| 3回目 | 8秒(4×2.0) |
4-5-2. CrawlerステートのRetry
Retry = [
{
ErrorEquals = ["Glue.CrawlerRunningException", ...]
IntervalSeconds = 30 # Crawlerは時間がかかるので長めに
MaxAttempts = 3
BackoffRate = 1.5
}
]
Crawlerが既に実行中の場合にリトライします。
4-5-3. NotifyErrorステート
NotifyError = {
Type = "Task"
Resource = "arn:aws:states:::sns:publish"
Parameters = {
TopicArn = aws_sns_topic.alert.arn
Subject = "【エラー】データパイプライン失敗"
Message = { ... }
}
End = true
}
| 設定 | 説明 |
|---|---|
| TopicArn | 通知先のSNSトピック |
| Subject | メールの件名 |
| Message | 通知内容(エラー詳細) |
4-6. フロー図(改善後)
5. 動作確認
5-1. 計画(プレビュー)
cd pipeiac02/tf
terraform plan
期待される出力:
Plan: 2 to add, 1 to change, 0 to destroy.
5-2. 適用(作成)
terraform apply
5-3. SNSサブスクリプション設定(手動)
メール通知を受け取るには、AWSコンソールで設定:
- SNS → トピック →
dp-alert - サブスクリプションの作成
- プロトコル: Email
- エンドポイント: あなたのメールアドレス
- 確認メールのリンクをクリック
5-4. 正常系テスト
リポジトリに含まれている test-data/error-test-ok.json を使用します。
# 正常なデータをアップロード
aws s3 cp test-data/error-test-ok.json s3://dp-raw-$(terraform output -raw bucket_suffix)/input/
# 実行状態を確認
SFN_ARN=$(aws stepfunctions list-state-machines --query 'stateMachines[?contains(name, `dp-pipeline`)].stateMachineArn' --output text)
sleep 10
aws stepfunctions list-executions --state-machine-arn $SFN_ARN --max-results 1 --query 'executions[0].{Status:status}'
期待される出力:
{
"Status": "SUCCEEDED"
}
5-5. エラー系テスト(意図的に失敗させる)
リポジトリに含まれている test-data/invalid-json.json を使用します。
# 不正なJSONをアップロード
aws s3 cp test-data/invalid-json.json s3://dp-raw-$(terraform output -raw bucket_suffix)/input/error-test.json
# 実行状態を確認(リトライ後に失敗するまで待つ)
sleep 30
aws stepfunctions list-executions --state-machine-arn $SFN_ARN --max-results 1 --query 'executions[0].{Status:status}'
期待される出力:
{
"Status": "FAILED"
}
→ SNSでメール通知が届く(設定済みの場合)
5-6. AWSコンソールで実行履歴を確認
- Step Functions → ステートマシン →
dp-pipeline - 失敗した実行をクリック
- グラフビューでエラー箇所を確認
5-7. 確認チェックリスト
| 確認項目 | 期待値 |
|---|---|
| SNSトピック |
dp-alert が存在 |
| 正常データ | SUCCEEDED |
| 異常データ | FAILED + SNS通知 |
| リトライ | 実行履歴で確認可能 |
6. まとめ
6-1. この記事でやったこと
| 項目 | 内容 |
|---|---|
| 理解 | Retry/Catch パターンを学んだ |
| 実装 | Step Functionsにエラー処理を追加 |
| 通知 | SNSでエラー通知を設定 |
6-2. 比喩の振り返り
| レストラン | AWS | 今回やったこと |
|---|---|---|
| 作り直し | Retry | 最大3回リトライ |
| 店長報告 | SNS通知 | エラー時にメール |
| 対応マニュアル | Catch | エラー処理ルール |
6-3. エラーハンドリングの全体像
6-4. 次回予告
第9回: CloudWatch監視 では、モニタリングを実装していきます。
- メトリクスの収集
- アラームの設定
- ダッシュボードの作成
レストランで言うと「厨房の監視カメラと警報システム」を作っていきます。