サーバーレスなデータワークフロー

AWS Kinesis と Lambda で学ぶストリーミングデータ

Maksim Pecherskiy

Data Engineer

前回の復習

z8HRk.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

現在のパーティション構成

FMktI.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

スピーダーをフィルタリング

IPwGK.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

日次集計

tvBXx.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

recordReaderS3 を開く

aIahM.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

recordReaderS3 を編集

...
import pytz
from datetime import datetime
tz = pytz.timezone('America/Los_Angeles')
AWS Kinesis と Lambda で学ぶストリーミングデータ

recordReaderS3 を編集

def record_created_handler(event, context):

... ## 45 の制限超過の最高速度を取得 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD :]
## オブジェクトキーを生成 fdate = datetime.now(tz).strftime("%Y%m%d/%HH%MM%SS") obj_key = f"speeders/{fdate}.csv" # speeders フォルダ内のファイル名
## オブジェクトを S3 に書き込み s3.put_object(Bucket='sd-vehicle-data', Key=obj_key, Body=too_fast.to_csv(sep=" ", index=False) )
AWS Kinesis と Lambda で学ぶストリーミングデータ

日時別にスピーダーを書き出す

s3QY4.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

speederAggregator

tvBXx.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

speederAggregator を作成

3pGYX.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

AWS Data Wrangler レイヤーを追加

7TO0V.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

speederAggregator のリソース

sO0op.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

定時トリガーを作成

4Qi4i.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

y0FkQ.png

cron(50 6 * * ? *)
  • 50分に
  • UTC 6時台
  • 毎月の毎日
  • 毎月
  • 毎週の毎日
  • 毎年
AWS Kinesis と Lambda で学ぶストリーミングデータ

speederAggregator コールバック

import boto3, pytz, pandas as pd

s3 = boto3.client("s3" ...)
...

def speeder_aggregator(event, context): tz = pytz.timezone('America/Los_Angeles') filter_date = datetime.now(tz).strftime("%Y%m%d")
AWS Kinesis と Lambda で学ぶストリーミングデータ

speederAggregator コールバック

def speeder_aggregator(event, context):
  ...
  objects = s3.list_objects_v2(
      Bucket='sd-vehicle-data', Prefix=f'speeders/{filter_date}')
  day_data = []

for obj in objects['Contents']: print(obj['Key']) day_record = s3.get_object(Bucket='sd-vehicle-data', Key = obj['Key']) day_data.append( pd.read_csv(day_record['Body'], delimiter = " ")) # 新規レコードを単一のデータフレームに結合 data = pd.concat(day_data) data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
AWS Kinesis と Lambda で学ぶストリーミングデータ

awswrangler パッケージ

import awswrangler as wr

session = boto3.Session(aws_access_key_id = AWS_KEY, aws_secret_access_key = AWS_SECRET, region_name="us-east-1")
df = wr.s3.read_csv(f"s3://sd-vehicle-data/speeders/{filter_date}") df.head()
AWS Kinesis と Lambda で学ぶストリーミングデータ

従来方式 vs awswrangler

def speeder_aggregator(event, context):
  ...
  objects = s3.list_objects_v2(
      Bucket='sd-vehicle-data', Prefix=f'speeders/{filter_date}')
  day_data = []

  for obj in objects['Contents']:
      print(obj['Key'])
      day_record = s3.get_object(Bucket='sd-vehicle-data', Key = obj['Key'])
      day_data.append(
          pd.read_csv(day_record['Body'], delimiter = " "))

  # 新規レコードを単一のデータフレームに結合
  data = pd.concat(day_data)
  data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
AWS Kinesis と Lambda で学ぶストリーミングデータ

従来方式 vs awswrangler

import awswrangler as wr

session = boto3.Session(aws_access_key_id = AWS_KEY, 
                        aws_secret_access_key = AWS_SECRET, region_name="us-east-1")
def speeder_aggregator(event, context):
    df = wr.s3.read_csv(f"s3://sd-vehicle-data/speeders/{filter_date}")      
    df.head()
AWS Kinesis と Lambda で学ぶストリーミングデータ

集計スピーダーファイルを書き出す

wr.s3.to_csv(df, 
             f"s3://sd-vehicle-data/speeders_daily/{filter_date}.csv", 
             boto3_session=session, 
             sep=" ", index=False)
AWS Kinesis と Lambda で学ぶストリーミングデータ

保存してテスト

agF2b.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

トリガーを有効化

COfYY.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

復習

u1Jwm.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

復習

7ZYPV.png

AWS Kinesis と Lambda で学ぶストリーミングデータ

復習

  • スケジューリング用 Cron
  • wr.s3.read_csv()
  • wr.s3.write_csv()
AWS Kinesis と Lambda で学ぶストリーミングデータ

Ayo berlatih!

AWS Kinesis と Lambda で学ぶストリーミングデータ

Preparing Video For Download...