첫 라이브 람다!

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

Maksim Pecherskiy

Data Engineer

지난 강의

StkGt.png

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

이번 강의

Zb4Va.png

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

공통 로그 형식

bc75da5f-1bf6-444c-80ad-49c180e1b8de 23:16:06.000 3FTEX1G5XAK844393 -76.6990172 2.481207 40
ff8e7131-408d-463b-8d07-d016419b0656 20:26:44.000 2LAXX1C8XAK844292 114.39239199999999 36.097577 90
f29a5b3d-d0fa-43c0-9e1a-e2a5cdb8be7a 8:10:47.000 3FTEX1G5XAK844393 108.58068100000001 34.79925 37
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

S3 파일 읽기

import pandas as pd

results in:

Response:
{
  "errorMessage": "Unable to import module 'lambda_function': No module named 'pandas'",
  "errorType": "Runtime.ImportModuleError"
}
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

핸들러 업데이트

#lambda_function.py
import json, boto3, pandas as pd
# 클라이언트 초기화
...

SPEED_ALERT_THRESHOLD = 45 ALERT_PHONE_NUMBER = "+1234567890"
# 작성된 레코드에서 데이터프레임을 가져오는 헬퍼 함수 def get_new_data(event): pd.read_csv()... ... return data
# Lambda 핸들러 def record_created_handler(event, context): data = get_new_data(event) ... sns.publish() ...
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

get_new_data()

def get_new_data(event):
    # 새 객체 키를 저장할 리스트 생성
    written_objects = []

    # 각 S3 이벤트 레코드 반복
    for record in event['Records']:

        # 확인할 변수 가져오기
        event_name = record['eventName']
        bucket_name = record['s3']['bucket']['name']
        obj_key = record['s3']['object']['key']
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

get_new_data()

def get_new_data(event):
        ...
        # sd-vehicle-data 버킷에서 생성된 이벤트인지 확인
        if event_name == 'ObjectCreated:Put' and bucket_name == 'sd-vehicle-data':

obj = s3.get_object(Bucket=bucket_name, Key = obj_key) df = pd.read_csv(obj['Body'], delimiter = " ", names=["record_id", "timestamp", "vin", "lon", "lat", "speed"])
written_objects.append(df)
# 새 레코드를 하나의 데이터프레임으로 결합 return pd.concat(written_objects)
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

record_created_handler()

SPEED_ALERT_THRESHOLD = 45
ALERT_PHONE_NUMBER = "+1234567890"
...
def record_created_handler(event, context):
    # 헬퍼 메서드 호출
    data = get_new_data(event)

## 최고 속도 계산 top_speeds = data.groupby(['vin'])['speed'].max().reset_index()
## 45 초과 속도 필터링 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD, :]
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

record_created_handler()

SPEED_ALERT_THRESHOLD = 45
ALERT_PHONE_NUMBER = "+1234567890"
...

def record_created_handler(event, context):
    ...
    ## SMS 전송
    sns.publish(PhoneNumber=ALERT_PHONE_NUMBER,
        Message="Speeding Alert \n" + too_fast.to_string())

## 아직 반환값은 사용하지 않지만, 무언가는 반환해야 함 totals = data.groupby(['vin'])['speed'].max().reset_index() return totals.to_csv(sep=" ", index=False)
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

Lambda 함수 테스트

jBpMo.png

KodWB.png

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

환경 변수 추가

import os
os.environ.get("ENV_VARIABLE_NAME", "DEFAULT_VALUE")
import os
SPEED_ALERT_THRESHOLD = os.environ.get("SPEED_ALERT_THRESHOLD", 45)
ALERT_PHONE_NUMBER = os.environ.get("ALERT_PHONE_NUMBER", None)
...

def record_created_handler(event, context):
    ...
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

환경 변수 추가

7jD6x.png

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

환경 변수 추가

eDRTX.png

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

yKjWC.png

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

트리거 추가

gLmbZD.png

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

트리거 추가

7F08p.png

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

트리거 추가

bPUtX.png

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

복습

ZAbfT.png

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

복습

m9jlv.png

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

복습

GTlyD.png

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

복습

XddOW.png

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

Lass uns üben!

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

Preparing Video For Download...