📁 完全なコードはGitHubで公開:GitHub: pipeiac02
シリーズ一覧(全12回)
Phase 1: 基盤構築
Phase 2: ワークフロー
Phase 3: セキュリティ・運用
Phase 4: 開発効率化
1. はじめに
1-1. 今回のゴール
Lambda関数を作成して、JSONからParquetへの変換処理を実装します。
| ゴール | 内容 |
|---|---|
| 設計 | ETL処理の流れを理解する |
| 実装 | TerraformでLambdaを作成する |
| 確認 | 変換処理が動くことを確認する |
1-2. 前回の振り返り
前回(第2回)では、S3バケットを作成しました。
| 作成済み | 今回作る |
|---|---|
| S3 Raw | Lambda ETL |
| S3 Processed | - |
1-3. この記事で作るもの
2. 比喩で理解する
2-1. Lambdaを「料理人」で考える
前回作った「冷蔵庫」と「盛り付け台」の間で、料理人が調理を担当します。
2-2. 比喩の図解(レストラン)
2-3. 料理人の仕事
| 手順 | 内容 |
|---|---|
| 1. 食材を取り出す | 冷蔵庫から食材を取る |
| 2. 調理する | 切る・焼く・味付け |
| 3. 盛り付ける | 皿に盛って完成 |
2-4. 技術の図解(AWS)
2-5. Lambdaの仕事
| 手順 | 内容 |
|---|---|
| 1. データを取り出す | S3 Rawからファイルを読む |
| 2. 変換する | JSON → Parquet |
| 3. 保存する | S3 Processedに書き込む |
2-6. 対応関係
| レストラン | AWS | 役割 |
|---|---|---|
| 冷蔵庫 | S3 Raw | 生の食材(JSON) |
| 料理人 | Lambda | 調理(変換処理) |
| 盛り付け台 | S3 Processed | 完成品(Parquet) |
| レシピ | Pythonコード | 調理手順 |
| 調理器具 | Pandas Layer | 変換ツール |
2-7. なぜ「サーバーレス」か
| 方式 | 例え | メリット |
|---|---|---|
| サーバー常駐 | 専属シェフ | 常にスタンバイ |
| Lambda | 派遣シェフ | 使った分だけ課金 |
3. Lambda設計の考え方
3-1. 処理の流れ
3-2. 入出力の整理
| 項目 | 内容 |
|---|---|
| 入力 | S3 Raw の JSONファイル |
| 出力 | S3 Processed の Parquetファイル |
| トリガー | Step Functions から呼び出し(後の回で設定) |
3-3. 必要なもの
| 要素 | 役割 | 料理で例えると |
|---|---|---|
| Lambda関数 | 処理本体 | 料理人 |
| Pythonコード | 変換ロジック | レシピ |
| Pandas Layer | データ処理ライブラリ | 調理器具 |
| IAMロール | 権限 | 調理場への入室許可 |
3-4. ファイルパスの設計
3-4-1. 入力(Raw)
s3://dp-raw-XXXXXX/input/ec-sales.json
3-4-2. 出力(Processed)
s3://dp-processed-XXXXXX/processed/year=2026/month=01/data_20260115_123456.parquet
3-5. なぜパーティション構造?
| 理由 | 説明 |
|---|---|
| クエリ効率 | 必要な月だけスキャン |
| コスト削減 | スキャン量が減る |
| 管理しやすい | 月ごとに整理される |
3-6. Lambda設定
| 設定 | 値 | 理由 |
|---|---|---|
| ランタイム | Python 3.12 | 最新安定版 |
| メモリ | 256 MB | Pandas処理に十分 |
| タイムアウト | 60秒 | 変換処理の余裕を持たせる |
| Layer | Pandas | データ処理に必要 |
4. Parquet形式とは
4-1. なぜJSONからParquetに変換するのか?
分析に適した形式に変えるためです。
4-2. 比喩で理解する
4-3. 料理で例えると
| 形式 | 例え |
|---|---|
| JSON | 買い物リスト(人が読みやすい) |
| Parquet | 冷蔵庫の整理棚(すぐ取り出せる) |
4-4. 行指向 vs 列指向
4-5. 分析時の違い
「priceの合計を出したい」場合:
| 形式 | 動作 | 効率 |
|---|---|---|
| JSON(行指向) | 全行を読んでpriceを抽出 | 遅い・コスト高 |
| Parquet(列指向) | price列だけ読む | 速い・コスト低 |
4-6. 比較まとめ
| 項目 | JSON | Parquet |
|---|---|---|
| 人の読みやすさ | ◎ | × |
| 機械の処理効率 | △ | ◎ |
| ファイルサイズ | 大きい | 小さい(圧縮) |
| 分析クエリ速度 | 遅い | 速い |
| 用途 | データ交換 | データ分析 |
4-7. このパイプラインでの使い分け
| 段階 | 形式 | 理由 |
|---|---|---|
| 入力 | JSON | 外部システムが出力しやすい |
| 分析 | Parquet | Athenaが高速にクエリできる |
5. 実装
5-1. ファイル構成
今回作成するファイル:
pipeiac02/
├── test-data/
│ └── ec-sales-03.json # ★テスト用サンプルデータ
└── tf/
├── lambda.tf # ★今回作成
├── iam.tf # ★今回作成(Lambda用ロール)
└── lambda/
└── etl.py # ★今回作成(Pythonコード)
5-2. Pythonコード(lambda/etl.py)
まず、変換処理のロジックを作成します。
コードの骨格:
# lambda/etl.py
import json
import boto3
import pandas as pd
from datetime import datetime
def handler(event, context):
# 1. イベントから入力キーを取得
# 2. S3からJSONを読み込み
# 3. Parquetに変換
# 4. S3に保存
# 5. 結果を返す
フルコード:
# lambda/etl.py
# JSON→Parquet変換を行うLambda関数
import json
import boto3
import pandas as pd
from datetime import datetime
import os
import logging
# ロガー設定
logger = logging.getLogger()
logger.setLevel(logging.INFO)
# S3クライアント
s3 = boto3.client('s3')
def handler(event, context):
"""
JSON → Parquet 変換処理
Args:
event: Step Functionsから渡されるイベント
- input_key: S3上のJSONファイルパス
context: Lambda実行コンテキスト
Returns:
dict: 処理結果
- statusCode: HTTPステータスコード
- output_key: 出力したParquetファイルのパス
- record_count: 処理したレコード数
"""
logger.info(f"Event: {json.dumps(event)}")
# 環境変数からバケット名を取得
raw_bucket = os.environ['RAW_BUCKET']
processed_bucket = os.environ['PROCESSED_BUCKET']
try:
# 1. 入力キーを取得
input_key = event['input_key']
logger.info(f"Processing: s3://{raw_bucket}/{input_key}")
# 2. S3からJSONを読み込み
response = s3.get_object(Bucket=raw_bucket, Key=input_key)
data = json.loads(response['Body'].read().decode('utf-8'))
# 3. DataFrameに変換
df = pd.DataFrame(data)
logger.info(f"Loaded {len(df)} records")
# 4-1. パーティションキーを生成(年/月でパーティション)
now = datetime.now()
year, month = now.strftime('%Y'), now.strftime('%m')
timestamp = now.strftime('%Y%m%d_%H%M%S')
# 4-2. 出力キーを生成
output_key = f"processed/year={year}/month={month}/data_{timestamp}.parquet"
# 5. Parquetに変換してS3に保存
parquet_buffer = df.to_parquet(index=False)
s3.put_object(Bucket=processed_bucket, Key=output_key, Body=parquet_buffer)
logger.info(f"Saved to: s3://{processed_bucket}/{output_key}")
# 6. 結果を返す
return {
'statusCode': 200,
'output_key': output_key,
'record_count': len(df)
}
except Exception as e:
logger.error(f"Error: {str(e)}")
raise
5-3. コード解説
| 処理 | コード | 料理で例えると |
|---|---|---|
| 材料を取る | s3.get_object() |
冷蔵庫から食材を出す |
| 下ごしらえ | pd.DataFrame(data) |
食材を切る |
| 調理 | df.to_parquet() |
焼く・味付け |
| 盛り付け | s3.put_object() |
皿に盛って配膳台へ |
5-4. Lambda ETL用IAM(iam.tf)
LambdaがS3にアクセスするための権限を設定します。
コードの骨格:
# iam.tf
# Lambda実行ロール
resource "aws_iam_role" "lambda_etl" {
# Lambdaが引き受けるロール
}
# S3アクセス権限
resource "aws_iam_role_policy" "lambda_s3" {
# Raw: 読み取り
# Processed: 書き込み
}
フルコード:
# iam.tf
# ================================
# Lambda ETL用IAM
# ================================
# Lambda実行ロール
# LambdaがS3にアクセスするための権限を付与
resource "aws_iam_role" "lambda_etl" {
name = "${var.project}-lambda-etl-role"
# Lambda用の信頼ポリシー
assume_role_policy = jsonencode({
Version = "2012-10-17"
Statement = [{
Action = "sts:AssumeRole"
Effect = "Allow"
Principal = { Service = "lambda.amazonaws.com" }
}]
})
tags = merge(local.common_tags, {
Name = "${var.project}-lambda-etl-role"
})
}
# Lambda用S3アクセスポリシー(最小権限)
# Rawバケットからの読み込み、Processedバケットへの書き込みのみ許可
resource "aws_iam_policy" "lambda_s3" {
name = "${var.project}-lambda-s3-policy"
description = "Lambda ETL用S3アクセスポリシー"
policy = jsonencode({
Version = "2012-10-17"
Statement = [
{
Sid = "ReadRawBucket"
Effect = "Allow"
Action = ["s3:GetObject"]
Resource = ["${aws_s3_bucket.raw.arn}/*"]
},
{
Sid = "WriteProcessedBucket"
Effect = "Allow"
Action = ["s3:PutObject"]
Resource = ["${aws_s3_bucket.processed.arn}/*"]
}
]
})
tags = local.common_tags
}
# Lambda S3ポリシーをアタッチ
resource "aws_iam_role_policy_attachment" "lambda_s3" {
role = aws_iam_role.lambda_etl.name
policy_arn = aws_iam_policy.lambda_s3.arn
}
# CloudWatch Logs出力用マネージドポリシーをアタッチ
resource "aws_iam_role_policy_attachment" "lambda_basic" {
role = aws_iam_role.lambda_etl.name
policy_arn = "arn:aws:iam::aws:policy/service-role/AWSLambdaBasicExecutionRole"
}
5-5. 権限の考え方
| バケット | 権限 | 理由 |
|---|---|---|
| Raw | 読み取りのみ | データを取得するだけ |
| Processed | 書き込みのみ | 変換結果を保存するだけ |
5-6. Lambda関数(lambda.tf)
コードの骨格:
# ソースコードをZIPに圧縮
data "archive_file" "lambda_etl" {
source_file = "${path.module}/lambda/etl.py"
output_path = "${path.module}/lambda/etl.zip"
}
# Lambda関数
resource "aws_lambda_function" "etl" {
function_name = "${var.project}-etl"
role = aws_iam_role.lambda_etl.arn
handler = "etl.handler"
runtime = "python3.12"
filename = data.archive_file.lambda_etl.output_path
layers = ["arn:aws:lambda:...:layer:AWSSDKPandas-Python312:14"] # Pandas Layer
environment { variables = { RAW_BUCKET = "...", PROCESSED_BUCKET = "..." } }
}
フルコード:
# lambda.tf
# ================================
# Lambda関数
# ================================
# Lambda関数のソースコードをZIPに圧縮
data "archive_file" "lambda_etl" {
type = "zip"
source_file = "${path.module}/lambda/etl.py"
output_path = "${path.module}/lambda/etl.zip"
}
# Lambda関数(ETL処理)
# JSON→Parquet変換を行うメイン処理
resource "aws_lambda_function" "etl" {
function_name = "${var.project}-etl"
role = aws_iam_role.lambda_etl.arn
handler = "etl.handler"
runtime = "python3.12"
# ZIPファイルのパスとハッシュ
filename = data.archive_file.lambda_etl.output_path
source_code_hash = data.archive_file.lambda_etl.output_base64sha256
# メモリとタイムアウト設定
memory_size = 256
timeout = 60
# AWS提供のPandas Layer(Parquet変換に必要)
layers = [
"arn:aws:lambda:${var.aws_region}:336392948345:layer:AWSSDKPandas-Python312:14"
]
# 環境変数でバケット名を渡す
environment {
variables = {
RAW_BUCKET = aws_s3_bucket.raw.id
PROCESSED_BUCKET = aws_s3_bucket.processed.id
}
}
tags = merge(local.common_tags, {
Name = "${var.project}-etl"
})
}
5-7. コード解説
| 設定 | 意味 | 料理で例えると |
|---|---|---|
handler |
実行する関数 | レシピの開始ページ |
runtime |
Python 3.12 | 調理場の設備 |
memory_size |
256 MB | 作業台の広さ |
timeout |
60秒 | 制限時間 |
layers |
Pandas(ARN直接指定) | 調理器具セット |
environment |
バケット名 | 材料置き場の場所 |
6. 動作確認
6-1. ディレクトリ構成の確認
cd pipeiac02/tf
ls -la lambda/
期待される出力:
etl.py
6-2. 計画(プレビュー)
terraform plan
期待される出力:
Plan: 4 to add, 0 to change, 0 to destroy.
追加されるリソース:
- IAMロール
- IAMポリシーアタッチメント
- IAMロールポリシー
- Lambda関数
6-3. 適用(作成)
terraform apply
Enter a value: が出たら yes を入力。
期待される出力:
Apply complete! Resources: 4 added, 0 changed, 0 destroyed.
6-4. Lambda関数の確認
aws lambda get-function --function-name dp-etl --query 'Configuration.{Name:FunctionName,Runtime:Runtime,Memory:MemorySize,Timeout:Timeout}'
期待される出力:
{
"Name": "dp-etl",
"Runtime": "python3.12",
"Memory": 256,
"Timeout": 60
}
6-5. テストデータの準備(サンプルJSON)
リポジトリに含まれている test-data/ec-sales-03.json を使用します。
# S3にアップロード
aws s3 cp test-data/ec-sales-03.json s3://dp-raw-$(terraform output -raw bucket_suffix)/input/ec-sales.json
6-6. Lambda手動実行
aws lambda invoke \
--function-name dp-etl \
--payload '{"input_key": "input/ec-sales.json"}' \
--cli-binary-format raw-in-base64-out \
response.json
cat response.json
期待される出力:
{
"statusCode": 200,
"body": {
"input": "s3://dp-raw-.../input/ec-sales.json",
"output": "s3://dp-processed-.../processed/year=2026/month=01/data_....parquet"
}
}
6-7. 出力ファイルの確認
aws s3 ls s3://dp-processed-$(terraform output -raw bucket_suffix)/processed/ --recursive
期待される出力:
2026-01-15 12:00:00 1234 processed/year=2026/month=01/data_20260115_120000.parquet
6-8. 確認チェックリスト
| 確認項目 | 期待値 |
|---|---|
| Lambda関数 |
dp-etl が存在 |
| ランタイム | Python 3.12 |
| メモリ | 256 MB |
| 実行結果 | statusCode: 200 |
| 出力ファイル | Parquetが生成されている |
7. まとめ
7-1. この記事でやったこと
| 項目 | 内容 |
|---|---|
| 設計 | Lambda ETLの処理フローを理解 |
| Parquet | なぜ変換するかを理解 |
| 実装 | TerraformでLambda関数を作成 |
| 確認 | JSON→Parquet変換を確認 |
7-2. 比喩の振り返り
| レストラン | AWS | 今回やったこと |
|---|---|---|
| 冷蔵庫 | S3 Raw | ✅ 前回作成 |
| 料理人 | Lambda | ✅ 今回作成 |
| 盛り付け台 | S3 Processed | ✅ 前回作成 |
| レシピ | Pythonコード | ✅ 今回作成 |
| 調理器具 | Pandas Layer | ✅ 今回追加 |
7-3. 作成したリソース
7-4. 処理の流れ
7-5. 次回予告
第4回: Glue Crawler では、メタデータ管理を実装していきます。
- Crawlerでスキーマ自動検出
- Glue Databaseの作成
- Athenaとの連携準備
レストランで言うと「メニュー表の自動更新」を作っていきます。