ユースケースおよびアーキパターン
ユースケース
Section titled “ユースケース”1. データレイクETL
Section titled “1. データレイクETL”シナリオ: - S3のログデータをParquet変換 - Athenaで分析可能に - 日次バッチ処理
実装: 1. Crawler でスキーマ検出 2. ETL Job で変換(CSV→Parquet) 3. パーティション作成 4. Athenaで分析2. データウェアハウス統合
Section titled “2. データウェアハウス統合”シナリオ: - 複数RDSからRedshiftへ集約 - データクレンジング - 増分ロード
実装: ソース: 複数RDS ETL Job: - データ抽出 - 変換・クレンジング - 重複削除 ターゲット: Redshift3. リアルタイムストリーミング
Section titled “3. リアルタイムストリーミング”シナリオ: - Kinesis Data Streamsからデータ処理 - リアルタイム変換 - S3/DynamoDB書き込み
実装: Kinesis Data Streams ↓ Glue Streaming Job ↓ S3(Parquet)/ DynamoDBアーキテクチャパターン
Section titled “アーキテクチャパターン”パターン1: シンプルETL
Section titled “パターン1: シンプルETL”構成: S3(ソースデータ): - CSV / JSON - 生データ
Glue Crawler: - スキーマ検出 - Data Catalog更新
Glue ETL Job: - データ読み込み - 変換処理 - Parquet出力
S3(変換後): - Parquet - パーティション化
Athena: - SQL分析
ETLスクリプト例(PySpark): - ApplyMapping(列マッピング) - Filter(不要データ除外) - Write(Parquet出力)パターン2: 増分処理
Section titled “パターン2: 増分処理”構成: S3(日次データ): - 新規ファイル追加
EventBridge: - スケジュール実行(日次)
Glue ETL Job: - ジョブブックマーク有効 - 新規データのみ処理 - 変換・出力
S3(蓄積データ): - Parquet - パーティション: year/month/day
メリット: - 処理時間短縮 - コスト削減 - 重複処理なしパターン3: 複数ソース統合
Section titled “パターン3: 複数ソース統合”構成: データソース: - RDS(顧客データ) - S3(ログデータ) - DynamoDB(トランザクション)
Glue ETL Job: - 各ソース読み込み - JOIN処理 - データクレンジング - 集計
Redshift: - データウェアハウス - 分析用
QuickSight: - ダッシュボード
変換処理: 1. 各ソースからデータ抽出 2. スキーマ正規化 3. JOINキーでマージ 4. 集計・変換 5. Redshift COPYパターン4: データ品質管理
Section titled “パターン4: データ品質管理”構成: S3(生データ) ↓ Glue DataBrew: - データプロファイリング - 異常検知 - クレンジングレシピ作成 ↓ Glue DataBrew Job: - レシピ適用 - データクレンジング ↓ S3(クレンジング済み) ↓ Glue ETL Job: - 変換処理 ↓ データウェアハウス
品質チェック: - NULL値検出 - データ型不整合 - 重複レコード - 統計的外れ値Crawler設定
Section titled “Crawler設定”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時実行)ETL Job(PySpark)
Section titled “ETL Job(PySpark)”import sysfrom awsglue.transforms import *from awsglue.utils import getResolvedOptionsfrom pyspark.context import SparkContextfrom awsglue.context import GlueContextfrom awsglue.job import Jobfrom awsglue.dynamicframe import DynamicFrame
# 引数取得args = getResolvedOptions(sys.argv, ['JOB_NAME'])
# コンテキスト初期化sc = SparkContext()glueContext = GlueContext(sc)spark = glueContext.spark_sessionjob = 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()ジョブブックマーク活用
Section titled “ジョブブックマーク活用”# ジョブブックマーク有効化datasource0 = glueContext.create_dynamic_frame.from_catalog( database="logs_database", table_name="raw_logs", transformation_ctx="datasource0", additional_options={ "jobBookmarkKeys": ["timestamp"], "jobBookmarkKeysSortOrder": "asc" })
# 増分データのみ処理される# 前回実行時の最終timestampから新規データのみ読み込みストリーミングETL
Section titled “ストリーミングETL”from awsglue.context import GlueContextfrom 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統合でセキュリティ強化