서버리스 데이터 워크플로

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

Maksim Pecherskiy

Data Engineer

돌아보기

되돌아보기

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

현재 파티셔닝

현재 파티셔닝

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

위반 차량 필터링

위반 차량 필터링

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

일별 집계

일별 집계

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

recordReaderS3 열기

recordReaderS3 열기

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로 배우는 스트리밍 데이터

날짜·시간별 위반 차량 저장

S3에 날짜·시간별로 저장

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

speederAggregator

speederAggregator

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

speederAggregator 생성

speederAggregator 만들기

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

AWS data wrangler 레이어 추가

AWS data wrangler 레이어 추가

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

speederAggregator 리소스

speederAggregator 리소스

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

시간 기반 트리거 생성

시간 기반 트리거 생성

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

Cron 식

cron(50 6 * * ? *)
  • 50분에
  • 6시(UTC)
  • 매달 모든 일자
  • 매월
  • 매주 모든 요일
  • 매년
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로 배우는 스트리밍 데이터

저장하고 테스트하기

저장 및 테스트

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

트리거 활성화

트리거 활성화

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

복습

복습

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

복습

복습

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

복습

  • 예약에는 Cron 사용
  • wr.s3.read_csv()
  • wr.s3.write_csv()
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

연습해 봅시다!

AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

Preparing Video For Download...