Skip to content

ユースケースおよびアーキパターン

シナリオ:
- S3のログデータをParquet変換
- Athenaで分析可能に
- 日次バッチ処理
実装:
1. Crawler でスキーマ検出
2. ETL Job で変換(CSV→Parquet)
3. パーティション作成
4. Athenaで分析
シナリオ:
- 複数RDSからRedshiftへ集約
- データクレンジング
- 増分ロード
実装:
ソース: 複数RDS
ETL Job:
- データ抽出
- 変換・クレンジング
- 重複削除
ターゲット: Redshift

3. リアルタイムストリーミング

Section titled “3. リアルタイムストリーミング”
シナリオ:
- Kinesis Data Streamsからデータ処理
- リアルタイム変換
- S3/DynamoDB書き込み
実装:
Kinesis Data Streams
Glue Streaming Job
S3(Parquet)/ DynamoDB
構成:
S3(ソースデータ):
- CSV / JSON
- 生データ
Glue Crawler:
- スキーマ検出
- Data Catalog更新
Glue ETL Job:
- データ読み込み
- 変換処理
- Parquet出力
S3(変換後):
- Parquet
- パーティション化
Athena:
- SQL分析
ETLスクリプト例(PySpark):
- ApplyMapping(列マッピング)
- Filter(不要データ除外)
- Write(Parquet出力)
構成:
S3(日次データ):
- 新規ファイル追加
EventBridge:
- スケジュール実行(日次)
Glue ETL Job:
- ジョブブックマーク有効
- 新規データのみ処理
- 変換・出力
S3(蓄積データ):
- Parquet
- パーティション: year/month/day
メリット:
- 処理時間短縮
- コスト削減
- 重複処理なし
構成:
データソース:
- RDS(顧客データ)
- S3(ログデータ)
- DynamoDB(トランザクション)
Glue ETL Job:
- 各ソース読み込み
- JOIN処理
- データクレンジング
- 集計
Redshift:
- データウェアハウス
- 分析用
QuickSight:
- ダッシュボード
変換処理:
1. 各ソースからデータ抽出
2. スキーマ正規化
3. JOINキーでマージ
4. 集計・変換
5. Redshift COPY
構成:
S3(生データ)
Glue DataBrew:
- データプロファイリング
- 異常検知
- クレンジングレシピ作成
Glue DataBrew Job:
- レシピ適用
- データクレンジング
S3(クレンジング済み)
Glue ETL Job:
- 変換処理
データウェアハウス
品質チェック:
- NULL値検出
- データ型不整合
- 重複レコード
- 統計的外れ値
import boto3
glue = boto3.client('glue')
# Crawler作成
response = glue.create_crawler(
Name='s3-logs-crawler',
Role='AWSGlueServiceRole-Crawler',
DatabaseName='logs_database',
Targets={
'S3Targets': [
{
'Path': 's3://my-bucket/logs/',
'Exclusions': [
'**.tmp',
'**_$folder$'
]
}
]
},
SchemaChangePolicy={
'UpdateBehavior': 'UPDATE_IN_DATABASE',
'DeleteBehavior': 'LOG'
},
Configuration='{"Version":1.0,"CrawlerOutput":{"Partitions":{"AddOrUpdateBehavior":"InheritFromTable"}}}',
Schedule='cron(0 0 * * ? *)' # 毎日0時実行
)
import sys
from awsglue.transforms import *
from awsglue.utils import getResolvedOptions
from pyspark.context import SparkContext
from awsglue.context import GlueContext
from awsglue.job import Job
from awsglue.dynamicframe import DynamicFrame
# 引数取得
args = getResolvedOptions(sys.argv, ['JOB_NAME'])
# コンテキスト初期化
sc = SparkContext()
glueContext = GlueContext(sc)
spark = glueContext.spark_session
job = Job(glueContext)
job.init(args['JOB_NAME'], args)
# ソースデータ読み込み
datasource0 = glueContext.create_dynamic_frame.from_catalog(
database="logs_database",
table_name="raw_logs",
transformation_ctx="datasource0"
)
# 列マッピング
applymapping1 = ApplyMapping.apply(
frame=datasource0,
mappings=[
("timestamp", "string", "event_time", "timestamp"),
("user_id", "string", "user_id", "string"),
("event_type", "string", "event_type", "string"),
("properties", "string", "properties", "string")
],
transformation_ctx="applymapping1"
)
# フィルタリング
filtered = Filter.apply(
frame=applymapping1,
f=lambda row: row["event_type"] in ["purchase", "signup", "login"]
)
# カスタム変換
def add_partition_keys(rec):
from datetime import datetime
dt = datetime.fromisoformat(rec["event_time"])
rec["year"] = str(dt.year)
rec["month"] = f"{dt.month:02d}"
rec["day"] = f"{dt.day:02d}"
return rec
mapped_dyf = Map.apply(
frame=filtered,
f=add_partition_keys
)
# Parquet出力
glueContext.write_dynamic_frame.from_options(
frame=mapped_dyf,
connection_type="s3",
connection_options={
"path": "s3://my-bucket/processed/",
"partitionKeys": ["year", "month", "day"]
},
format="parquet",
format_options={
"compression": "snappy"
},
transformation_ctx="datasink"
)
job.commit()
# ジョブブックマーク有効化
datasource0 = glueContext.create_dynamic_frame.from_catalog(
database="logs_database",
table_name="raw_logs",
transformation_ctx="datasource0",
additional_options={
"jobBookmarkKeys": ["timestamp"],
"jobBookmarkKeysSortOrder": "asc"
}
)
# 増分データのみ処理される
# 前回実行時の最終timestampから新規データのみ読み込み
from awsglue.context import GlueContext
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
glueContext = GlueContext(spark.sparkContext)
# Kinesis Data Streamsから読み込み
kinesis_options = {
"streamARN": "arn:aws:kinesis:ap-northeast-1:123456789012:stream/my-stream",
"startingPosition": "TRIM_HORIZON",
"inferSchema": "true",
"classification": "json"
}
streaming_df = glueContext.create_data_frame.from_options(
connection_type="kinesis",
connection_options=kinesis_options,
transformation_ctx="streaming_df"
)
# ウィンドウ集計
from pyspark.sql.functions import window, count, avg
windowed = streaming_df \
.withWatermark("timestamp", "10 minutes") \
.groupBy(
window("timestamp", "5 minutes"),
"user_id"
) \
.agg(
count("*").alias("event_count"),
avg("value").alias("avg_value")
)
# S3へストリーミング出力
query = windowed.writeStream \
.format("parquet") \
.option("path", "s3://my-bucket/streaming-output/") \
.option("checkpointLocation", "s3://my-bucket/checkpoints/") \
.partitionBy("window") \
.start()
query.awaitTermination()

アーキテクチャ設計のポイント:

  • Crawler で自動スキーマ管理
  • ジョブブックマークで増分処理
  • Parquet変換でコスト削減
  • パーティション設計で効率化
  • DataBrew でデータ品質向上
  • Step Functions でオーケストレーション
  • VPC統合でセキュリティ強化