Amazon EventBridge + AWS Step Functions の統合ワークショップでシンプルなものが公式になかったので作りました。
Amazon EventBridge とは
Amazon EventBridge は、AWS のサービスや自作アプリケーション、外部 SaaS から発生する「イベント」を受け取り、あらかじめ定めたルールに従って適切な宛先(ターゲット)へ振り分ける、サーバーレスのイベントバスです。サーバーの管理は不要で、使った分だけ課金されます。
中心となる考え方は「イベント駆動(Event-Driven)」です。たとえば「S3 にファイルが置かれた」「EC2 の状態が変わった」といった出来事がイベントとして流れてきて、EventBridge はその内容をイベントパターンで判定し、条件に合致したものだけを Lambda や Step Functions などのターゲットに渡します。イベントを出す側(プロデューサー)と受け取る側(コンシューマー)が直接つながらず、EventBridge を挟んで疎結合になるため、あとから処理を追加・変更しやすいのが利点です。
今回のハンズオンでは、S3 の input/ にオブジェクトが置かれたことを EventBridge が検知し、Step Functions を起動する「起点」として使います。なお、定期実行(cron のようなスケジュール起動)をしたい場合は EventBridge Scheduler という機能も用意されています。
補足として、EventBridge の配信は「少なくとも1回(at-least-once)」であり、まれに同じイベントが重複して届くことがあります。そのため、起動されるバッチ処理は「同じイベントが2回来ても結果が壊れない(冪等)」ように作っておくのが基本です。
AWS Step Functions とは
AWS Step Functions は、複数の処理を「ワークフロー」としてつなぎ、実行順序・条件分岐・リトライ・エラー処理をまとめて管理してくれる、サーバーレスのオーケストレーションサービスです。
ワークフローは「ステートマシン」と呼ばれ、その定義を Amazon States Language(ASL)という JSON 形式で記述します。ステートマシンは複数の「ステート(状態)」の組み合わせでできており、代表的なものに、処理を実行する Task、条件で分岐する Choice、並列実行する Parallel、繰り返す Map、待機する Wait などがあります。処理の流れをコンソール上で図として可視化できるため、全体像を把握しやすいのも特徴です。
Step Functions の強みは、各処理を「起動して終わりを待ち、成功なら次へ、失敗なら別の経路へ」と制御できる点にあります。リトライ(Retry)やエラー捕捉(Catch)を宣言的に書けるので、複雑なバッチの信頼性を高められます。とくに Glue・Athena・EMR などとの .sync(Run a Job)統合を使うと、外部ジョブの完了までステートが待機し、その成否で後続を分岐できます。
今回のハンズオンでは、EventBridge から起動され、Lambda での加工・Choice での種別分岐・Athena .sync での集計・エラー処理・結果の書き込みまでを、一本のワークフローとして束ねる「司令塔」の役割を担います。
なお、ワークフローには課金体系と上限の異なる Standard と Express の2種類があり、長時間・低頻度なら Standard、短時間・高頻度なら Express が向いています。
.sync(Run a Job)と通常呼び出しの違い
Step Functions から他の AWS サービスを呼び出すとき、その「待ち方」には主に2つのモードがあります。ジョブを起動してすぐ次へ進む Request Response(通常の呼び出し)と、ジョブの完了まで待ち受ける .sync(Run a Job)です。ASL 上では、呼び出す Resource の末尾に .sync を付けるかどうかで切り替わります。
通常呼び出し(Request Response) これは「起動命令を出して、受け付けられたら即座に次のステートへ進む」動き方です。API を呼び出して「受理しました」という応答が返ってきた時点で、そのステートは成功として完了扱いになります。ジョブが実際に終わったかどうかは待ちません。
たとえば Glue ジョブを Request Response で起動すると、Step Functions は「ジョブの開始を受け付けた」ことだけを確認して次に進みます。ジョブ本体はバックグラウンドで走り続けますが、Step Functions はその完了も成否も把握しません。投げっぱなしでよい処理(fire-and-forget)には向いていますが、後続のステートが「まだ処理中のデータ」を掴んでしまう危険があります。
.sync(Run a Job) これは「ジョブを起動し、その完了までステートを待機させ、成功/失敗を受け取ってから次へ進む」動き方です。ジョブが終わるまでそのステートはアクティブなまま待ち続け、ジョブが成功すれば次のステートへ、失敗すれば Catch へ、と成否に応じて制御できます。
Step Functions は裏側で対象サービスの完了イベントを(EventBridge 経由などで)検知して待機を解除するため、開発者が自前でポーリングのループを書く必要がありません。Glue の glue:startJobRun.sync、Athena の athena:startQueryExecution.sync、EMR Serverless の emr-serverless:startJobRun.sync などが代表例です。

さっそくやってみる Phase1
1. 環境準備
ではIAMユーザーを作成しシークレットキーとアクセスキーを発行します。 権限は、Amazon S3、Amazon EventBridge、AWS StepFunctions、AWS Lambda、Amazon Athena、AWS Glue Schema Registry へのフルアクセスをも持っているものとします。

設定が終わったらAWS CLIが使える環境で aws configure を設定してクレデンシャルをセットしておきます。
2. 必要な環境変数の設定
以下を実行しておきます。
export AWS_REGION=ap-northeast-1
export ACCOUNT_ID=$(aws sts get-caller-identity --query Account --output text)
export BUCKET=batch-handson-${ACCOUNT_ID} # 全世界で一意な名前にする
export FUNC=batch-handson-transform
export SM=batch-handson-sm
export RULE=batch-handson-rule3. S3 バケット作成 + EventBridge 連携を有効化
S3 は既定では EventBridge にイベントを送りません。明示的に有効化する作業が必要になります。
aws s3api create-bucket --bucket $BUCKET --region $AWS_REGION \
--create-bucket-configuration LocationConstraint=$AWS_REGION
aws s3api put-bucket-notification-configuration --bucket $BUCKET \
--notification-configuration '{"EventBridgeConfiguration":{}}'4. Lambda 関数用IAMロール
S3バケットからオブジェクトを読み取り、CloudWatchLogsに処理結果を入れますので必要な権限を保有したIAMロールを作成します。
cat > lambda-trust.json <<'EOF'
{"Version":"2012-10-17","Statement":[{"Effect":"Allow",
"Principal":{"Service":"lambda.amazonaws.com"},"Action":"sts:AssumeRole"}]}
EOF
export LAMBDA_ROLE_ARN=$(aws iam create-role \
--role-name ${FUNC}-role \
--assume-role-policy-document file://lambda-trust.json \
--query Role.Arn --output text)
# CloudWatch Logs への書き込み
aws iam attach-role-policy --role-name ${FUNC}-role \
--policy-arn arn:aws:iam::aws:policy/service-role/AWSLambdaBasicExecutionRole
# 対象バケットの読み書きだけを許可(最小権限)
cat > lambda-s3.json <<EOF
{"Version":"2012-10-17","Statement":[{"Effect":"Allow",
"Action":["s3:GetObject","s3:PutObject"],
"Resource":"arn:aws:s3:::${BUCKET}/*"}]}
EOF
aws iam put-role-policy --role-name ${FUNC}-role \
--policy-name s3-rw --policy-document file://lambda-s3.json5. Lambda 関数の作成
EventBridgeが S3バケットの /input フォルダにオブジェクトがアップロードされたというイベントをもとにStep Functionsを起動します。 Step Functionsは Lambda関数を起動し /inputフォルダからオブジェクトを読み取ります。その後処理を行い /outputフォルダに保存します。
cat > lambda_function.py <<'EOF'
import boto3, urllib.parse
s3 = boto3.client("s3")
def handler(event, context):
detail = event["detail"]
bucket = detail["bucket"]["name"]
key = urllib.parse.unquote_plus(detail["object"]["key"])
body = s3.get_object(Bucket=bucket, Key=key)["Body"].read().decode("utf-8")
transformed = body.upper() # ←「修正」の例
filename = key.split("/")[-1]
out_key = f"output/upper-{filename}" # 別名で output/ へ
s3.put_object(Bucket=bucket, Key=out_key,
Body=transformed.encode("utf-8"))
return {"input": key, "output": out_key, "bytes": len(transformed)}
EOF
zip function.zip lambda_function.py
# ロール作成直後は反映待ちが必要なことがある
sleep 10
export LAMBDA_ARN=$(aws lambda create-function \
--function-name $FUNC \
--runtime python3.12 \
--role $LAMBDA_ROLE_ARN \
--handler lambda_function.handler \
--timeout 30 \
--zip-file fileb://function.zip \
--query FunctionArn --output text)
処理はシンプルです。S3 に置かれたテキストファイルを読み取り → 中身を全部大文字に変換 → output/ に別名で書き戻します。
6. Step Functions 用 IAM ロール
StepFunctions が上記で作成したLambda関数を起動するための権限を保有したIAMロールを作成します。
cat > sfn-trust.json <<'EOF'
{"Version":"2012-10-17","Statement":[{"Effect":"Allow",
"Principal":{"Service":"states.amazonaws.com"},"Action":"sts:AssumeRole"}]}
EOF
export SFN_ROLE_ARN=$(aws iam create-role \
--role-name ${SM}-role \
--assume-role-policy-document file://sfn-trust.json \
--query Role.Arn --output text)
cat > sfn-invoke.json <<EOF
{"Version":"2012-10-17","Statement":[{"Effect":"Allow",
"Action":"lambda:InvokeFunction","Resource":"${LAMBDA_ARN}"}]}
EOF
aws iam put-role-policy --role-name ${SM}-role \
--policy-name invoke-lambda --policy-document file://sfn-invoke.json
7. ステートマシンを作成
StepFunctionsでは1個1個のジョブ設定を ステートマシン と呼びます。
cat > definition.json <<EOF
{
"Comment": "Phase1: read -> transform -> write",
"StartAt": "ProcessObject",
"States": {
"ProcessObject": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "${LAMBDA_ARN}", "Payload.\$": "\$" },
"End": true
}
}
}
EOF
sleep 10
export SM_ARN=$(aws stepfunctions create-state-machine \
--name $SM \
--definition file://definition.json \
--role-arn $SFN_ROLE_ARN \
--type STANDARD \
--query stateMachineArn --output text)
"Payload.$": "$" は「EventBridge から受け取った入力(S3 イベント全体)をそのまま Lambda に渡す」といことを意味していて、とてもシンプルなものです。コンソールではこの様に見えます。

8. EventBridge 用 IAM ロール
EventBridge が上記で作成したステートマシンを起動できるようにするIAMロールを作成します。
cat > eb-trust.json <<'EOF'
{"Version":"2012-10-17","Statement":[{"Effect":"Allow",
"Principal":{"Service":"events.amazonaws.com"},"Action":"sts:AssumeRole"}]}
EOF
export EB_ROLE_ARN=$(aws iam create-role \
--role-name ${RULE}-role \
--assume-role-policy-document file://eb-trust.json \
--query Role.Arn --output text)
cat > eb-start.json <<EOF
{"Version":"2012-10-17","Statement":[{"Effect":"Allow",
"Action":"states:StartExecution","Resource":"${SM_ARN}"}]}
EOF
aws iam put-role-policy --role-name ${RULE}-role \
--policy-name start-sfn --policy-document file://eb-start.json9. EventBridge ルール作成 + ターゲット紐付け
ルールとはS3イベントのどういう動作に合致すればイベントを処理するか?という設定です。以下の場合 /inputにオブジェクトがputされればイベント発火、となります。 ターゲットとはルールに合致した場合何を起動するか?を指定します。以下の例では上記で作成したステートマシンを起動します。
cat > pattern.json <<EOF
{
"source": ["aws.s3"],
"detail-type": ["Object Created"],
"detail": {
"bucket": { "name": ["${BUCKET}"] },
"object": { "key": [ { "prefix": "input/" } ] }
}
}
EOF
aws events put-rule --name $RULE \
--event-pattern file://pattern.json --state ENABLED
sleep 10
cat > targets.json <<EOF
[{ "Id": "sfn-target", "Arn": "${SM_ARN}", "RoleArn": "${EB_ROLE_ARN}" }]
EOF
aws events put-targets --rule $RULE --targets file://targets.json10. テスト
ではいよいよテストです。 S3バケットにオブジェクトをputすると、EventBridge → StepFunctions → Lambda → S3 と処理が流れます。 以下の様な内容が出力されれば完成です。
--------------------------------------------------------------------------------------------
| ListExecutions |
+-----------------------------------------------------------------------------+------------+
| name | status |
+-----------------------------------------------------------------------------+------------+
| 9a9040e5-7410-651e-4e44-67a2eb51a050_7deccdab-40ce-7dd4-1b04-82cab24d8a03 | SUCCEEDED |
+-----------------------------------------------------------------------------+------------+
HELLO HANDSON
Step Functions のコンソールでは以下の様に表示されています。

CloudWatchのロググループ /aws/lambda/batch-handson-transform にも以下の通りログが出力されています

StepFunctions の処理分岐
ではここからPhase2です。 Phase 2 では、Phase 1 の「読み取って加工して書き込む」という一本道の処理に、Step Functions ならではの制御——条件分岐とエラー処理—を加えます。ファイルの種類に応じて処理を振り分け、一時的な失敗は自動で再試行し、それでもうまくいかないものは安全に隔離する、という信頼性の高いバッチの形を目指します。全体は「種別判定」「分岐」「加工(再試行・失敗捕捉つき)」「エラー退避」の4つのパートで構成されます。

種別判定(Classify) ワークフローの入り口では、まず専用の Lambda がファイルの種類を判定します。EventBridge から渡ってきた S3 イベントを受け取り、オブジェクトのキー(ファイル名)から拡張子を調べて、txt・csv・その他(other)のいずれかに分類します。この段階では実際の加工は行わず、後続のステートが判断に使えるよう、対象のバケット名・キー・判定した種別(fileType)だけをまとめて返すのが役割です。処理の「仕分け札」を作る工程だと考えると分かりやすいです。
Choice で分岐(RouteByType) 次に Choice ステートが、Classify の返した fileType を見て進む先を決めます。txt または csv であれば加工処理(Transform)へ、それ以外であればエラー退避(HandleError)へと振り分けます。ここで重要なのが Default の指定です。想定した条件のどれにも当てはまらなかったイベントは Default の経路に流れます。Default を用意しておかないと、未知の拡張子が来たときにワークフロー自体がエラーで停止してしまうため、「その他はすべてエラー退避へ」という受け皿として必ず設定します。
加工(Transform)— Retry と Catch 分岐を通過した txt / csv は、加工用の Lambda で処理されます。txt なら本文を大文字化、csv なら内容を整形し、結果を output/ に別名で書き込みます。Phase 1 との最大の違いは、この加工ステートに Retry と Catch という2段構えの安全装置を付けている点です。
Retry は「一時的な失敗を自動でやり直す」仕組みです。Lambda の一時的なスロットリングや瞬間的なサービス側の不調など、時間をおけば成功する類のエラーに対して、間隔を空けながら(指数バックオフで 2秒→4秒→8秒のように)決められた回数まで自動で再試行します。人手を介さずに一過性の障害を吸収するのが狙いです。
Catch は「再試行しても最終的に失敗したときの受け皿」です。Retry で規定回数を試してもなお失敗した場合や、そもそも再試行で回復しない種類の失敗(今回の教材では、本文に特定の文字列が含まれていたらわざと例外を投げる仕掛けで再現します)に対して、ワークフローを異常終了させるのではなく、エラー退避の経路へ流します。このとき、元の入力情報(バケット名やキー)を捨てずにエラー内容を追記して渡すよう設定しておくのがポイントで、これにより退避先の処理が「どのファイルが・なぜ失敗したか」を把握できます。
エラー退避(HandleError) Choice の Default から来た「対応していない種別のファイル」と、Transform の Catch から来た「加工に失敗したファイル」は、いずれもこのエラー退避ステートに集約されます。ここでは問題のあったファイルを error/ プレフィックスへコピーして隔離し、理由を記録します。正常な出力(output/)と失敗したもの(error/)が物理的に分かれるため、後から「何が処理され、何が弾かれたか」を一目で確認できます。分岐由来の失敗も加工由来の失敗も、経路は違えど最終的に同じ退避先へ合流する設計になっています。
さっそくやってみる Phase2
では作っていきます。
1. ファイル種別判別用Lambda関数の作成
拡張子を見て txt / csv / other を返すだけの軽い関数です。IAM ロールは Phase 1 の Lambda ロールを使い回します。
export CLASSIFY=batch-handson-classify
export ERRORFN=batch-handson-error
cat > classify.py <<'EOF'
import urllib.parse
def handler(event, context):
detail = event["detail"]
bucket = detail["bucket"]["name"]
key = urllib.parse.unquote_plus(detail["object"]["key"])
ext = key.rsplit(".", 1)[-1].lower() if "." in key else ""
file_type = ext if ext in ("txt", "csv") else "other"
return {"bucket": bucket, "key": key, "fileType": file_type}
EOF
zip classify.zip classify.py
export CLASSIFY_ARN=$(aws lambda create-function \
--function-name $CLASSIFY \
--runtime python3.12 --role $LAMBDA_ROLE_ARN \
--handler classify.handler --timeout 10 \
--zip-file fileb://classify.zip \
--query FunctionArn --output text)2. Error Lambda(error/ へ退避)を作成
問題のあったファイルを error/ プレフィックスにコピーして隔離する処理を行います。
cat > errorfn.py <<'EOF'
import boto3
s3 = boto3.client("s3")
def handler(event, context):
bucket = event["bucket"]
key = event["key"]
reason = event.get("error", {}).get("Error", "UnsupportedFileType")
filename = key.split("/")[-1]
dest = f"error/{filename}"
s3.copy_object(Bucket=bucket, Key=dest,
CopySource={"Bucket": bucket, "Key": key})
print(f"quarantined {key} -> {dest} (reason={reason})")
return {"quarantined": dest, "reason": reason}
EOF
zip errorfn.zip errorfn.py
export ERROR_ARN=$(aws lambda create-function \
--function-name $ERRORFN \
--runtime python3.12 --role $LAMBDA_ROLE_ARN \
--handler errorfn.handler --timeout 10 \
--zip-file fileb://errorfn.zip \
--query FunctionArn --output text)3. Transform Lambda(既存)を改修
次にPhase1で作成したLambda関数を修正します。 Phase1ではこのLambda関数はEventBridgeからの意ンプトっと受け取っていましたが、 Phase 2 では Classify が整形した {bucket, key, fileType} を受け取る形に変えます。あわせて、Catch を試すための 「本文に FAIL が含まれていたらわざと例外を投げる」 仕掛けと、csv の簡単な加工を追加します。
cat > lambda_function.py <<'EOF'
import boto3, csv, io
s3 = boto3.client("s3")
def handler(event, context):
bucket, key, ftype = event["bucket"], event["key"], event["fileType"]
body = s3.get_object(Bucket=bucket, Key=key)["Body"].read().decode("utf-8")
if "FAIL" in body: # ← Catch 動作確認用のわざと失敗
raise RuntimeError("intentional failure for testing Catch")
filename = key.split("/")[-1]
if ftype == "txt":
out = body.upper()
out_key = f"output/upper-{filename}"
else: # csv
rows = list(csv.reader(io.StringIO(body)))
out = f"# rows={len(rows)}\n" + body # 行数を先頭に付けるだけの例
out_key = f"output/counted-{filename}"
s3.put_object(Bucket=bucket, Key=out_key, Body=out.encode("utf-8"))
return {"input": key, "output": out_key, "fileType": ftype}
EOF
zip function.zip lambda_function.py
aws lambda update-function-code \
--function-name $FUNC --zip-file fileb://function.zip4. Step Functions ロールに新しく作成する Lambda関数 の起動権限を追加
Phase 1 では Transform しか呼べませんでした。3つとも呼べるよう、命名規則に合わせたワイルドカードで上書きします。
cat > sfn-invoke.json <<EOF
{"Version":"2012-10-17","Statement":[{"Effect":"Allow",
"Action":"lambda:InvokeFunction",
"Resource":"arn:aws:lambda:${AWS_REGION}:${ACCOUNT_ID}:function:batch-handson-*"}]}
EOF
aws iam put-role-policy --role-name ${SM}-role \
--policy-name invoke-lambda --policy-document file://sfn-invoke.json5. ステートマシンの更新
Phase1で作成したものを今回用に置き換えます。
cat > definition.json <<EOF
{
"Comment": "Phase2: classify -> choice -> transform (retry/catch) / error",
"StartAt": "Classify",
"States": {
"Classify": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "${CLASSIFY_ARN}", "Payload.\$": "\$" },
"OutputPath": "\$.Payload",
"Next": "RouteByType"
},
"RouteByType": {
"Type": "Choice",
"Choices": [
{ "Or": [
{ "Variable": "\$.fileType", "StringEquals": "txt" },
{ "Variable": "\$.fileType", "StringEquals": "csv" }
], "Next": "Transform" }
],
"Default": "HandleError"
},
"Transform": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "${FUNC}", "Payload.\$": "\$" },
"Retry": [ {
"ErrorEquals": ["Lambda.ServiceException","Lambda.TooManyRequestsException","States.TaskFailed"],
"IntervalSeconds": 2, "MaxAttempts": 3, "BackoffRate": 2.0
} ],
"Catch": [ {
"ErrorEquals": ["States.ALL"],
"ResultPath": "\$.error",
"Next": "HandleError"
} ],
"End": true
},
"HandleError": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "${ERROR_ARN}", "Payload.\$": "\$" },
"End": true
}
}
}
EOF
aws stepfunctions update-state-machine \
--state-machine-arn $SM_ARN --definition file://definition.json
6. テスト
では4パターンでテストします。
# ① txt → 正常(大文字化されて output/upper-*.txt)
echo "hello phase2" > a.txt
aws s3 cp a.txt s3://${BUCKET}/input/a.txt
# ② csv → 正常(行数付きで output/counted-*.csv)
printf "id,name\n1,alice\n2,bob\n" > b.csv
aws s3 cp b.csv s3://${BUCKET}/input/b.csv
# ③ その他拡張子 → Choice の Default で error/ へ退避
echo "unsupported" > c.log
aws s3 cp c.log s3://${BUCKET}/input/c.log
# ④ txt だが FAIL を含む → Retry 3回後に Catch で error/ へ退避
echo "this will FAIL" > d.txt
aws s3 cp d.txt s3://${BUCKET}/input/d.txt
数十秒だって処理が終わったかテストしてみます。
aws stepfunctions list-executions --state-machine-arn $SM_ARN \
--max-results 6 --query 'executions[].{name:name,status:status}' --output table
aws s3 ls s3://${BUCKET}/output/
aws s3 ls s3://${BUCKET}/error/
S3バケットに以下の処理済ファイルが展開されています。
ファイル | 元 | どう処理されたか |
|---|---|---|
| ① a.txt(txt) | Transform で大文字化 → |
| ② b.csv(csv) | Transform で行数付与 → |
| ③ c.log(その他) | Choice の Default で弾かれ → |
| ④ d.txt(FAIL入り) | Retry 3回後 Catch で捕捉 → |
ちなみに実行履歴の一番上(一番新しいもの)は上記の通り3回実行したあとエラー処理となっています。少しややこしいですが、エラー処理が成功しているためステータスは成功になっています。

コンソールで詳細を見ると以下の通りの実行履歴となっています。


.sync(Run a Job)
これまでの Lambda 呼び出しは、Step Functions がその場で関数を実行し、戻り値を受け取って次へ進む、同期的で完結した処理でした。一方 Athena のクエリは「投げてすぐ終わる」ものではなく、データ量によっては数秒から数分かかる非同期のジョブです。
startQueryExecution.sync を使うと、Step Functions はクエリを起動したあと、その完了を検知するまでステートを待機させます。クエリが成功すれば次のステートへ、失敗すれば Catch へ——という制御が、ポーリングのループを自前で書くことなく実現できます。Phase3では「起動する・完了を待つ・成否で分岐する」という一連の動きを、コードではなくワークフローの定義として宣言的に書ける点を、実際に動かして確かめます。

さっそくやってみる:Phase3
1. 変数の確認と追加
まず実行に必要なあたりを環境変数に登録します。
export ATHENA_DB=batch_handson_db
export ATHENA_TABLE=sales
export GLUE_DB_LOC=s3://${BUCKET}/glue-db/ # Glue DB のメタ用
export ATHENA_OUT=s3://${BUCKET}/athena-results/ # クエリ結果の出力先
export ATHENA_WG=batch-handson-wg # ワークグループ
echo "BUCKET=[$BUCKET] SM_ARN=[$SM_ARN] CLASSIFY_ARN=[$CLASSIFY_ARN] ERROR_ARN=[$ERROR_ARN] FUNC=[$FUNC]"
BUCKET=[batch-handson-917561075114] SM_ARN=[arn:aws:states:ap-northeast-1:917561075114:stateMachine:batch-handson-sm] CLASSIFY_ARN=[arn:aws:lambda:ap-northeast-1:917561075114:function:batch-handson-classify] ERROR_ARN=[arn:aws:lambda:ap-northeast-1:917561075114:function:batch-handson-error] FUNC=[batch-handson-transform]
2. Glue データベースとテーブルを作成
S3 上の csv を Athena が SQL で読めるよう、「表」として定義します。AthenaはGlueに登録しているデータカタログを自動でテーブル定義として使用しますのでこの2つはセットのサービスになります。 まずデータベースを作成します。
aws glue create-database \
--database-input "{\"Name\":\"${ATHENA_DB}\",\"LocationUri\":\"${GLUE_DB_LOC}\"}"
次にテーブルを作成します。id(整数)と name(文字列)の2列を持つ csv として、datasetフォルダ を参照するテーブルを定義します。
cat > table-input.json <<EOF
{
"Name": "${ATHENA_TABLE}",
"TableType": "EXTERNAL_TABLE",
"Parameters": { "classification": "csv", "skip.header.line.count": "1" },
"StorageDescriptor": {
"Columns": [
{ "Name": "id", "Type": "int" },
{ "Name": "name", "Type": "string" }
],
"Location": "s3://${BUCKET}/dataset/",
"InputFormat": "org.apache.hadoop.mapred.TextInputFormat",
"OutputFormat": "org.apache.hadoop.hive.ql.io.HiveIgnoreKeyTextOutputFormat",
"SerdeInfo": {
"SerializationLibrary": "org.apache.hadoop.hive.serde2.OpenCSVSerde",
"Parameters": { "separatorChar": ",", "quoteChar": "\"" }
}
}
}
EOF
aws glue create-table \
--database-name $ATHENA_DB \
--table-input file://table-input.json3. Athena ワークグループを作成
Athena は自動で上記で設定したデータベースを読み取ってクエリを発行します。上記コマンドではGlueの機能を使ってS3のdatasetバケットのcsvをテーブルとして登録しています。 これらの連携により、S3のdatasetバケットに登録されているCSVにAthenaがSQLで分析クエリを発行できるようになります。
Athenaはクエリの実行結果の保存を、指定されたS3バケットにおこないます。この設定はワークグループで行います。
aws athena create-work-group \
--name $ATHENA_WG \
--configuration "ResultConfiguration={OutputLocation=${ATHENA_OUT}}"4. Step Functions ロールに Athena / Glue / S3 権限を追加
Phase2まではステートマシンはLambda関数を起動していましたが、Phase3ではAthenaとGlueを操作します。また直接S3も操作するため権限が必要となります。
cat > sfn-athena.json <<EOF
{
"Version": "2012-10-17",
"Statement": [
{
"Sid": "AthenaQuery",
"Effect": "Allow",
"Action": [
"athena:startQueryExecution",
"athena:stopQueryExecution",
"athena:getQueryExecution",
"athena:getQueryResults"
],
"Resource": "arn:aws:athena:${AWS_REGION}:${ACCOUNT_ID}:workgroup/${ATHENA_WG}"
},
{
"Sid": "GlueCatalog",
"Effect": "Allow",
"Action": [ "glue:GetDatabase", "glue:GetTable", "glue:GetPartitions" ],
"Resource": [
"arn:aws:glue:${AWS_REGION}:${ACCOUNT_ID}:catalog",
"arn:aws:glue:${AWS_REGION}:${ACCOUNT_ID}:database/${ATHENA_DB}",
"arn:aws:glue:${AWS_REGION}:${ACCOUNT_ID}:table/${ATHENA_DB}/*"
]
},
{
"Sid": "S3DataAndResults",
"Effect": "Allow",
"Action": [ "s3:GetObject", "s3:PutObject", "s3:ListBucket", "s3:GetBucketLocation" ],
"Resource": [
"arn:aws:s3:::${BUCKET}",
"arn:aws:s3:::${BUCKET}/*"
]
}
]
}
EOF
aws iam put-role-policy --role-name ${SM}-role \
--policy-name athena-access --policy-document file://sfn-athena.json5. ステートマシンに Athena .sync ステートを組み込み
このPhase3では.syncでAthenaを起動します。これによりジョブの処理完了を確認するポーリング処理を書き込むことなく自動でStepFunctionsが制御してくれるようになります。
cat > definition.json <<EOF
{
"Comment": "Phase3: classify -> choice(txt/csv/other) -> transform / athena.sync / error",
"StartAt": "Classify",
"States": {
"Classify": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "${CLASSIFY_ARN}", "Payload.\$": "\$" },
"OutputPath": "\$.Payload",
"Next": "RouteByType"
},
"RouteByType": {
"Type": "Choice",
"Choices": [
{ "Variable": "\$.fileType", "StringEquals": "txt", "Next": "Transform" },
{ "Variable": "\$.fileType", "StringEquals": "csv", "Next": "AthenaAggregate" }
],
"Default": "HandleError"
},
"Transform": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "${FUNC}", "Payload.\$": "\$" },
"Retry": [ {
"ErrorEquals": ["Lambda.ServiceException","Lambda.TooManyRequestsException","States.TaskFailed"],
"IntervalSeconds": 2, "MaxAttempts": 3, "BackoffRate": 2.0
} ],
"Catch": [ { "ErrorEquals": ["States.ALL"], "ResultPath": "\$.error", "Next": "HandleError" } ],
"End": true
},
"AthenaAggregate": {
"Type": "Task",
"Resource": "arn:aws:states:::athena:startQueryExecution.sync",
"Parameters": {
"QueryString": "SELECT name, COUNT(*) AS cnt FROM ${ATHENA_DB}.${ATHENA_TABLE} GROUP BY name",
"WorkGroup": "${ATHENA_WG}",
"QueryExecutionContext": { "Database": "${ATHENA_DB}" }
},
"Retry": [ {
"ErrorEquals": ["States.TaskFailed"],
"IntervalSeconds": 5, "MaxAttempts": 2, "BackoffRate": 2.0
} ],
"Catch": [ { "ErrorEquals": ["States.ALL"], "ResultPath": "\$.error", "Next": "HandleError" } ],
"End": true
},
"HandleError": {
"Type": "Task",
"Resource": "arn:aws:states:::lambda:invoke",
"Parameters": { "FunctionName": "${ERROR_ARN}", "Payload.\$": "\$" },
"End": true
}
}
}
EOF
aws stepfunctions update-state-machine \
--state-machine-arn $SM_ARN --definition file://definition.jsonステートマシンは以下の様になっています。

6. テスト
ではいよいよテスト実行です。 まず、Athena が集計する実データを dataset フォルダに置きます。 datasetフォルダはinputフォルダとは異なりますのでこの手順ではEventBridgeはイベントをS3から受け取りません。
printf "id,name\n1,alice\n2,bob\n3,alice\n" > sales.csv
aws s3 cp sales.csv s3://${BUCKET}/dataset/sales.csv
次に inputフォルダにオブジェクトを配置することでEventBridgeにイベントを発火させ、ステートマシンを起動させます。
# ① csv → Athena .sync で集計が走る
printf "id,name\n1,alice\n2,bob\n3,alice\n" > trigger.csv
aws s3 cp trigger.csv s3://${BUCKET}/input/trigger.csv
# ② txt → Transform(大文字化)
echo "hello phase3" > e.txt
aws s3 cp e.txt s3://${BUCKET}/input/e.txt
# ③ その他 → Default で error/ へ
echo "nope" > f.log
aws s3 cp f.log s3://${BUCKET}/input/f.log
Athenaの処理には少し時間がかかります。ステートマシン側のコンソールを見てみると1件だけ実行中になっています。

て終了した後集計結果を確認します。
# 出力された結果ファイルの一覧
aws s3 ls s3://${BUCKET}/athena-results/
# 一番新しい .csv 結果ファイルの中身を表示(.metadata ファイルは除外)
LATEST=$(aws s3 ls s3://${BUCKET}/athena-results/ --recursive \
| awk '{print $4}' | grep '\.csv$' | sort | tail -1)
echo "result file: $LATEST"
aws s3 cp s3://${BUCKET}/${LATEST} -
2026-09-08 23:10:54 35 6164862d-ead7-4972-b66b-d40e2a08c6b6.csv
2026-09-08 23:10:54 104 6164862d-ead7-4972-b66b-d40e2a08c6b6.csv.metadata
result file: athena-results/6164862d-ead7-4972-b66b-d40e2a08c6b6.csv
"name","cnt"
"bob","1"
"alice","2"
集計前のCSVは以下ですのでaliceが集計されて2になっていることがわかります。
id,name
1,alice
2,bob
3,alice
