Firehose 전송 스트림 다루기

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

Maksim Pecherskiy

Data Engineer

스트림 생성 준비

KF414O.png

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

스트림 생성 준비

YgVbe.png

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

스트림 생성 준비

BBYB4.png

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

스트림 생성 준비

Ue73e.png

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

스트림 생성 준비

qUq0E.png

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

Role ARN 가져오기

qH99c.png

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

boto3 클라이언트 초기화

import boto3

firehose = boto3.client('firehose', 
                        aws_access_key_id=AWS_KEY_ID, 
                        aws_secret_access_key=AWS_SECRET, 
                        region_name='us-east-1')
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

스트림 생성하기

res = firehose.create_delivery_stream(

DeliveryStreamName = "gps-delivery-stream",
DeliveryStreamType = "DirectPut",
S3DestinationConfiguration = {
"RoleARN": "arn:aws:iam::0000000:role/firehoseDeliveryRole",
"BucketARN": "arn:aws:s3:::sd-vehicle-data"
}
)
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

스트림 생성 응답

print(res['DeliveryStreamARN'])
# 새 스트림의 ARN
"arn:aws:firehose:us-east-1:0000000:deliverystream/gps-delivery-stream"
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

스트림 준비 완료

z0Fp4.png

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

스트림에 쓰기

La0zZ.png

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

텔레매틱스 하드웨어

OGHUHC.png

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

텔레매틱스 데이터 전송

wCJz8.png

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

단일 레코드

{

'record_id': '939ed1d1-1740-420c-8906-445278573c7f', # <-- 고유 레코드 ID
'timestamp': '4:25:06.000', # <-- 측정 시각
'vin': '4FTEX4944AK844294', # <-- 차량 ID
'lon': 106.9447146, # <-- 차량 경도
'lat': -6.3385652, # <-- 차량 위도
'speed': 25 # <-- 차량 속도
}
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

레코드 유입 중

vRg6A.png

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

다른 사용 사례

gu6oT.png

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

다른 사용 사례

XFD7E.png

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

패턴

h7hUU.png

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

레코드 전송

res = firehose.put_record(

DeliveryStreamName='gps-delivery-stream',
Record = { 'Data': payload }
)
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

레코드 전송

    Record = {
        'Data': payload 
    }
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

레코드 전송

레코드 예시
record = {
 'record_id': '939ed1d1-1740-420c-8906-445278573c7f',
 'timestamp': '4:25:06.000','vin': '4FTEX4944AK844294',
 'lon': 106.9447146,'lat': -6.338565200000001,
 'speed': 25}
전송 형식(단일 문자열)
"939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 
4FTEX4944AK844294 106.9447146 -6.338565200000001 25"
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

레코드 전송

payload = " ".join(
    str(value) for value in record.values()
)
print(payload)
"939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 
4FTEX4944AK844294 106.9447146 -6.338565200000001 25"
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

종합하기

record = {
 'record_id': '939ed1d1-1740-420c-8906-445278573c7f',
 'timestamp': '4:25:06.000','vin': '4FTEX4944AK844294',
 'lon': 106.9447146,'lat': -6.338565200000001, 'speed': 25}

payload = " ".join( str(value) for value in record.values() )
#"939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 4FTEX4944AK844294 106.9447146 -6.338565200000001 25"
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

종합하기

res = firehose.put_record(
    DeliveryStreamName='gps-delivery-stream',
    Record = {
        'Data': payload + "\n" #<-- 줄바꿈!
    }
)
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

생성된 파일

McYNJ.png

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

샘플 데이터

939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 4FTEX4944AK844294 106.9447146 -6.338565200000001 25
f29a5b3d-d0fa-43c0-9e1a-e2a5cdb8be7a 8:10:47.000 3FTEX1G5XAK844393 108.58068100000001 34.79925 37
ff8e7131-408d-463b-8d07-d016419b0656 20:26:44.000 2LAXX1C8XAK844292 114.39239199999999 36.097577 90
bc75da5f-1bf6-444c-80ad-49c180e1b8de 23:16:06.000 3FTEX1G5XAK844393 -76.6990172 2.481207 40
7bdcf779-444e-4313-83da-140461933aeb 22:28:44.000 5FTEX1MAXAK844295 -47.0145295 -21.4649238 40
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

생성된 파일

X0px9.png

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

S3 클라이언트 생성

# boto3 S3 클라이언트 생성
s3 = boto3.client('s3', 
                  aws_access_key_id=AWS_KEY_ID, 
                  aws_secret_access_key=AWS_SECRET, 
                  region_name='us-east-1')
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

DataFrame으로 읽기

# S3에서 객체 가져오기
obj_data = s3.get_object(Bucket='sd-vehicle-data', Key=KEY_YOU_COPIED)
# 객체를 DataFrame으로 읽기
vehicle_data = pd.read_csv(
    data['Body'], 
    delimiter = " ", 
    names=["record_id", "timestamp", "vin", "lon", "lat", "speed"]))
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

vehicle_data

record_id timestamp vin lon lat speed
0 939ed1d1... 4:25:06.000 4FTEX4944AK844294 106.945 -6.33857 25
1 f29a5b3d... 8:10:47.000 3FTEX1G5XAK844393 108.581 34.7993 37
2 ff8e7131... 20:26:44.000 2LAXX1C8XAK844292 114.392 36.0976 90
3 bc75da5f... 23:16:06.000 3FTEX1G5XAK844393 -76.699 2.48121 40
4 7bdcf779... 22:28:44.000 5FTEX1MAXAK844295 -47.0145 -21.4649 40
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

복습

qUq0E.png

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

복습

res = firehose.create_delivery_stream(

DeliveryStreamName = "gps-delivery-stream",
DeliveryStreamType = "DirectPut",
S3DestinationConfiguration = {
"RoleARN": "arn:aws:iam::0000000:role/firehoseDeliveryRole",
"BucketARN": "arn:aws:s3:::sd-vehicle-data",
}
)
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

복습

wCJz8.png

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

복습

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

복습

record_id timestamp vin lon lat speed
0 939ed1d1... 4:25:06.000 4FTEX4944AK844294 106.945 -6.33857 25
1 f29a5b3d... 8:10:47.000 3FTEX1G5XAK844393 108.581 34.7993 37
2 ff8e7131... 20:26:44.000 2LAXX1C8XAK844292 114.392 36.0976 90
3 bc75da5f... 23:16:06.000 3FTEX1G5XAK844393 -76.699 2.48121 40
4 7bdcf779... 22:28:44.000 5FTEX1MAXAK844295 -47.0145 -21.4649 40
AWS Kinesis와 Lambda로 배우는 스트리밍 데이터

연습해 봅시다!

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

Preparing Video For Download...