import boto3
import json
import io
import pyarrow as pa
import pyarrow.parquet as pq
from concurrent.futures import ThreadPoolExecutor, as_completed
s3 = boto3.client('s3')
BUCKET = 'ad-clicks-raw-dev-262969866776'
PREFIX = 'raw-clicks/'
paginator = s3.get_paginator('list_objects_v2')
# 找出所有 ad_id/date 資料夾(同上)
ad_prefixes = []
for year_page in paginator.paginate(Bucket=BUCKET, Prefix=PREFIX, Delimiter='/'):
for prefix in year_page.get('CommonPrefixes', []):
year_prefix = prefix['Prefix']
for month_page in paginator.paginate(Bucket=BUCKET, Prefix=year_prefix, Delimiter='/'):
for month_prefix in month_page.get('CommonPrefixes', []):
for day_page in paginator.paginate(Bucket=BUCKET, Prefix=month_prefix['Prefix'], Delimiter='/'):
for day_prefix in day_page.get('CommonPrefixes', []):
for ad_page in paginator.paginate(Bucket=BUCKET, Prefix=day_prefix['Prefix'], Delimiter='/'):
for ad_prefix in ad_page.get('CommonPrefixes', []):
ad_prefixes.append(ad_prefix['Prefix'])
print(f"Found {len(ad_prefixes)} ad_id/date combinations")
def process_ad_prefix(ad_prefix):
s3_local = boto3.client('s3')
parts = ad_prefix.rstrip('/').split('/')
date = f"{parts[1]}/{parts[2]}/{parts[3]}"
ad_id = parts[4].replace('ad_id=', '')
records = []
pag = s3_local.get_paginator('list_objects_v2')
for page in pag.paginate(Bucket=BUCKET, Prefix=ad_prefix):
for obj in page.get('Contents', []):
key = obj['Key']
if not key.endswith('.json'):
continue
response = s3_local.get_object(Bucket=BUCKET, Key=key)
records.append(json.loads(response['Body'].read().decode('utf-8')))
if not records:
return f"Skip: {ad_prefix}"
table = pa.Table.from_pylist(records)
buffer = io.BytesIO()
pq.write_table(table, buffer)
buffer.seek(0)
output_key = f"parquet-clicks/{date}/ad_id={ad_id}/ad_clicks.parquet"
s3_local.put_object(Bucket=BUCKET, Key=output_key, Body=buffer.getvalue())
return f"Done: {output_key} ({len(records)} records)"
# 平行處理,同時跑 20 個
with ThreadPoolExecutor(max_workers=20) as executor:
futures = {executor.submit(process_ad_prefix, p): p for p in ad_prefixes}
for future in as_completed(futures):
print(future.result())
print("All completed!")
轉換後的 Parquet 檔案,按日期與 ad_id 分區存放
s3://ad-clicks-raw-dev-xxxxxxxx/
parquet-clicks/
2026/
01/
02/
ad_id=ad_1/
ad_clicks.parquet ← 這天 ad_1 的所有 click 事件
ad_id=ad_2/
ad_clicks.parquet ← 這天 ad_2 的所有 click 事件
...
完成後可先用線上的工具檢查 Parquet 檔,可以看到原本的 JSON 資料有明確的 Schema ,不同檔案的資料也被 merge 再一起。
Medium: 使用 Serverless 架構設計廣告點擊系統 — 用 AWS Glue Python Shell Job 把 S3 上的 JSON 檔轉成 Parquet 檔