iT邦幫忙

2026 iThome 鐵人賽

DAY 14
0
Software Development

使用 Serverless 架構設計廣告點擊系統 系列 第 14 篇

Day 14: 用 AWS Glue Python Shell Job 把 S3 上的 JSON 檔轉成 Parquet 檔(下)

  • 分享至 

  • xImage
  •  

AWS Glue Python Shell Job

Glue Python Shell Job script

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 檔


上一篇
Day 13: 用 AWS Glue Python Shell Job 把 S3 上的 JSON 檔轉成 Parquet 檔(中)
下一篇
Day 15: 用 EMR 將資料一次性遷移到 Aurora PostgreSQL 的完整紀錄(上)
系列文
使用 Serverless 架構設計廣告點擊系統 共 16 篇
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言