เวิร์กโฟลว์ข้อมูลแบบ Serverless

Streaming Data ด้วย AWS Kinesis และ Lambda

Maksim Pecherskiy

Data Engineer

ทบทวนที่ผ่านมา

z8HRk.png

Streaming Data ด้วย AWS Kinesis และ Lambda

การแบ่งพาร์ติชันในปัจจุบัน

FMktI.png

Streaming Data ด้วย AWS Kinesis และ Lambda

กรองยานพาหนะที่ขับเร็วเกินกำหนด

IPwGK.png

Streaming Data ด้วย AWS Kinesis และ Lambda

รวมข้อมูลตามวัน

tvBXx.png

Streaming Data ด้วย AWS Kinesis และ Lambda

เปิด recordReaderS3

aIahM.png

Streaming Data ด้วย AWS Kinesis และ Lambda

แก้ไข recordReaderS3

...
import pytz
from datetime import datetime
tz = pytz.timezone('America/Los_Angeles')
Streaming Data ด้วย AWS Kinesis และ Lambda

แก้ไข recordReaderS3

def record_created_handler(event, context):

... ## Get top speeds that exceed the limit of 45 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD :]
## Generate object key fdate = datetime.now(tz).strftime("%Y%m%d/%HH%MM%SS") obj_key = f"speeders/{fdate}.csv" # filename in speeders folder
## Write the object to S3 s3.put_object(Bucket='sd-vehicle-data', Key=obj_key, Body=too_fast.to_csv(sep=" ", index=False) )
Streaming Data ด้วย AWS Kinesis และ Lambda

บันทึกข้อมูลผู้ขับเร็วตามวันและเวลา

s3QY4.png

Streaming Data ด้วย AWS Kinesis และ Lambda

speederAggregator

tvBXx.png

Streaming Data ด้วย AWS Kinesis และ Lambda

สร้าง speederAggregator

3pGYX.png

Streaming Data ด้วย AWS Kinesis และ Lambda

เพิ่ม AWS data wrangler layer

7TO0V.png

Streaming Data ด้วย AWS Kinesis และ Lambda

ทรัพยากรของ speederAggregator

sO0op.png

Streaming Data ด้วย AWS Kinesis และ Lambda

สร้าง timed trigger

4Qi4i.png

Streaming Data ด้วย AWS Kinesis และ Lambda

y0FkQ.png

cron(50 6 * * ? *)
  • นาทีที่ 50
  • ของชั่วโมงที่ 6 (UTC)
  • ของทุกวันในเดือน
  • ของทุกเดือน
  • ของทุกวันในสัปดาห์
  • ของทุกปี
Streaming Data ด้วย AWS Kinesis และ Lambda

callback ของ 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")
Streaming Data ด้วย AWS Kinesis และ Lambda

callback ของ 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 = " ")) # Concatenate new records into a single dataframe. data = pd.concat(day_data) data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
Streaming Data ด้วย 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()
Streaming Data ด้วย 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 = " "))

  # Concatenate new records into a single dataframe.
  data = pd.concat(day_data)
  data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
Streaming Data ด้วย 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()
Streaming Data ด้วย AWS Kinesis และ Lambda

บันทึกไฟล์รวมข้อมูลผู้ขับเร็ว

wr.s3.to_csv(df, 
             f"s3://sd-vehicle-data/speeders_daily/{filter_date}.csv", 
             boto3_session=session, 
             sep=" ", index=False)
Streaming Data ด้วย AWS Kinesis และ Lambda

บันทึกและทดสอบ

agF2b.png

Streaming Data ด้วย AWS Kinesis และ Lambda

เปิดใช้งาน trigger

COfYY.png

Streaming Data ด้วย AWS Kinesis และ Lambda

ทบทวน

u1Jwm.png

Streaming Data ด้วย AWS Kinesis และ Lambda

ทบทวน

7ZYPV.png

Streaming Data ด้วย AWS Kinesis และ Lambda

ทบทวน

  • Cron สำหรับการกำหนดเวลา
  • wr.s3.read_csv()
  • wr.s3.write_csv()
Streaming Data ด้วย AWS Kinesis และ Lambda

มาฝึกกันเถอะ!

Streaming Data ด้วย AWS Kinesis และ Lambda

Preparing Video For Download...