변환용 Lambda

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

Maksim Pecherskiy

Data Engineer

1장에서는...

2020-07-15_06-18.png

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

2장에서는...

2020-07-15_06-20.png

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

이 장에서는...

2020-07-15_06-27.png

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

이전 접근법

2020-07-14_07-42.png

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

S3 일괄 처리 vs Lambda 변환

S3에서 일괄 처리

  • S3 객체 쓰기 시 Lambda 트리거 사용
  • 변환까지 지연이 김
  • 정제 전 원시 데이터 저장
  • Firehose 대상이 S3여야 함

Lambda 변환으로 처리

  • Firehose 중간에서 Lambda 트리거 사용
  • 즉시 변환됨
  • 정제된 데이터만 저장
  • Firehose의 다른 대상도 허용
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

유입 데이터

record_id timestamp vin lon lat speed
939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 4FTEX4944AK844294 106.9447146 -6.3385652 25
f29a5b3d-d0fa-43c0-9e1a-e2a5cdb8be7a 8:10:47.000 3FTEX1G5XAK844393 108.580681 34.79925 37
ff8e7131-408d-463b-8d07-d016419b0656 20:26:44.000 2LAXX1C8XAK844292 114.392392 36.097577 90
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

변환 Lambda

2020-07-15_06-38.png

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

샘플 이벤트

{
  "invocationId": "invocationIdExample",
  "deliveryStreamArn": "arn:aws:firehose:us-east-1:458912630:deliverystream/gps-delivery-stream",
  "region": "us-east-1",
  "records": [
    {

"recordId": "49546986683135544286507457936321625675700192471156785154",
"approximateArrivalTimestamp": 1495072949453,
"data": "NjQuMjQyLjg4LjEwIC0gLSBbMDcvTWFyLzIwMDQ6MTY6MTA6MDIgLTA4MDBdICJHRVQgL21haWxtYW4vbGlzdGluZm8vaHNkaXZpc2lvbiBIVFRQLzEuMSIgMjAwIDYyOTE=="
} ] }
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

Base64

2020-07-19_10-54.png

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

데이터 처리

import base64
from datetime import datetime as dt

def convert_timestamp(record_time_val):
# 오늘 날짜를 문자열로 가져옵니다 today = dt.today().strftime("%Y-%m-%d")
# 오늘 날짜와 레코드 시간을 결합해 datetime 객체를 만듭니다 new_ts = dt.strptime(f"{today} {record_time_val}", "%Y-%m-%d %H:%M:%S.%f")
# datetime 객체를 보기 좋게 포맷합니다 return new_ts.strftime("%Y-%m-%dT%H:%M:%S")
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

데이터 처리

def lambda_handler(event, context):
    output = []

for record in event['records']: payload = base64.b64decode(record['data'])
payload = payload.decode()
payload = payload.split(" ")
payload[1] = convert_timestamp(payload[1])
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

데이터 처리

    for record in event['records']:
        ...
        payload = " ".join(payload)
        payload_enc = base64.b64encode(payload.encode())

output.append({
'recordId': record['recordId'],
'result': 'Ok',
'data': payload_enc
})
return {'records': output}
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

복습

import base64
def lambda_handler(event, context):

output = [] for record in event['records']: # 레코드 반복
payload = base64.b64decode(record['data']).decode() # 페이로드 디코드
# 수정
payload_enc = base64.b64encode(payload.encode()) # 재인코드
output.append({ # 딕셔너리에 담기 'recordId': record['recordId'], 'result': 'Ok', 'data': payload_enc, })
return {'records': output}
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

AWS에서 Lambda 생성

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

Lambda 생성 스크립트

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

연습해 봅시다!

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

Preparing Video For Download...