Làm việc với luồng phân phối Firehose

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Maksim Pecherskiy

Data Engineer

Sẵn sàng tạo stream

KF414O.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Sẵn sàng tạo stream

YgVbe.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Sẵn sàng tạo stream

BBYB4.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Sẵn sàng tạo stream

Ue73e.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Sẵn sàng tạo stream

qUq0E.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Lấy Role ARN

qH99c.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Khởi tạo client 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')
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tạo stream!

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"
}
)
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Phản hồi tạo stream

print(res['DeliveryStreamARN'])
# ARN của stream mới
"arn:aws:firehose:us-east-1:0000000:deliverystream/gps-delivery-stream"
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Stream is ready

z0Fp4.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ghi vào stream

La0zZ.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Thiết bị telematics

OGHUHC.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Gửi dữ liệu telematics

wCJz8.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Một bản ghi

{

'record_id': '939ed1d1-1740-420c-8906-445278573c7f', # <-- ID bản ghi duy nhất
'timestamp': '4:25:06.000', # <-- thời điểm đo
'vin': '4FTEX4944AK844294', # <-- ID xe
'lon': 106.9447146, # <-- kinh độ vị trí xe
'lat': -6.3385652, # <-- vĩ độ vị trí xe
'speed': 25 # <-- tốc độ xe
}
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Bản ghi đổ về

vRg6A.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Trường hợp dùng khác

gu6oT.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Trường hợp dùng khác

XFD7E.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Mẫu hình

h7hUU.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Gửi một bản ghi

res = firehose.put_record(

DeliveryStreamName='gps-delivery-stream',
Record = { 'Data': payload }
)
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Gửi một bản ghi

    Record = {
        'Data': payload 
    }
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Gửi một bản ghi

Bản ghi của chúng ta trông như sau
record = {
 'record_id': '939ed1d1-1740-420c-8906-445278573c7f',
 'timestamp': '4:25:06.000','vin': '4FTEX4944AK844294',
 'lon': 106.9447146,'lat': -6.338565200000001,
 'speed': 25}
Chuỗi cần gửi (một dòng)
"939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 
4FTEX4944AK844294 106.9447146 -6.338565200000001 25"
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Gửi một bản ghi

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"
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ghép lại

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"
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ghép lại

res = firehose.put_record(
    DeliveryStreamName='gps-delivery-stream',
    Record = {
        'Data': payload + "\n" #<-- Xuống dòng!
    }
)
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tệp đã tạo

McYNJ.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Dữ liệu mẫu

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
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tệp đã tạo

X0px9.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tạo client S3

# Tạo client S3 của boto3.
s3 = boto3.client('s3', 
                  aws_access_key_id=AWS_KEY_ID, 
                  aws_secret_access_key=AWS_SECRET, 
                  region_name='us-east-1')
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Đọc dữ liệu vào DataFrame

# Lấy object từ S3
obj_data = s3.get_object(Bucket='sd-vehicle-data', Key=KEY_YOU_COPIED)
# Đọc object vào DataFrame
vehicle_data = pd.read_csv(
    data['Body'], 
    delimiter = " ", 
    names=["record_id", "timestamp", "vin", "lon", "lat", "speed"]))
Xử lý dữ liệu streaming với AWS Kinesis và 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
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ôn tập

qUq0E.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ôn tập

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",
}
)
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ôn tập

wCJz8.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ôn tập

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ôn tập

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
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ayo berlatih!

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Preparing Video For Download...