Travailler avec le flux de livraison Firehose

Streaming de données avec AWS Kinesis et Lambda

Maksim Pecherskiy

Data Engineer

Prêt à créer le flux

KF414O.png

Streaming de données avec AWS Kinesis et Lambda

Prêt à créer le flux

YgVbe.png

Streaming de données avec AWS Kinesis et Lambda

Prêt à créer le flux

BBYB4.png

Streaming de données avec AWS Kinesis et Lambda

Prêt à créer le flux

Ue73e.png

Streaming de données avec AWS Kinesis et Lambda

Prêt à créer le flux

qUq0E.png

Streaming de données avec AWS Kinesis et Lambda

Obtenir l'ARN du rôle

qH99c.png

Streaming de données avec AWS Kinesis et Lambda

Initialiser le 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')
Streaming de données avec AWS Kinesis et Lambda

Créer le flux

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"
}
)
Streaming de données avec AWS Kinesis et Lambda

Réponse de création du flux

print(res['DeliveryStreamARN'])
# ARN du nouveau flux
"arn:aws:firehose:us-east-1:0000000:deliverystream/gps-delivery-stream"
Streaming de données avec AWS Kinesis et Lambda

Flux prêt

z0Fp4.png

Streaming de données avec AWS Kinesis et Lambda

Écriture dans le flux

La0zZ.png

Streaming de données avec AWS Kinesis et Lambda

Matériel télématique

OGHUHC.png

Streaming de données avec AWS Kinesis et Lambda

Envoi des données télématiques

wCJz8.png

Streaming de données avec AWS Kinesis et Lambda

Un enregistrement

{

'record_id': '939ed1d1-1740-420c-8906-445278573c7f', # <-- ID d'enregistrement unique
'timestamp': '4:25:06.000', # <-- heure de mesure
'vin': '4FTEX4944AK844294', # <-- ID véhicule
'lon': 106.9447146, # <-- longitude du véhicule
'lat': -6.3385652, # <-- latitude du véhicule
'speed': 25 # <-- vitesse du véhicule
}
Streaming de données avec AWS Kinesis et Lambda

Arrivée des enregistrements

vRg6A.png

Streaming de données avec AWS Kinesis et Lambda

Autre cas d'usage

gu6oT.png

Streaming de données avec AWS Kinesis et Lambda

Autre cas d'usage

XFD7E.png

Streaming de données avec AWS Kinesis et Lambda

Motifs

h7hUU.png

Streaming de données avec AWS Kinesis et Lambda

Envoi d'un enregistrement

res = firehose.put_record(

DeliveryStreamName='gps-delivery-stream',
Record = { 'Data': payload }
)
Streaming de données avec AWS Kinesis et Lambda

Envoi d'un enregistrement

    Record = {
        'Data': payload 
    }
Streaming de données avec AWS Kinesis et Lambda

Envoi d'un enregistrement

À quoi ressemble notre enregistrement
record = {
 'record_id': '939ed1d1-1740-420c-8906-445278573c7f',
 'timestamp': '4:25:06.000','vin': '4FTEX4944AK844294',
 'lon': 106.9447146,'lat': -6.338565200000001,
 'speed': 25}
Ce que nous voulons envoyer (une chaîne)
"939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 
4FTEX4944AK844294 106.9447146 -6.338565200000001 25"
Streaming de données avec AWS Kinesis et Lambda

Envoi d'un enregistrement

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"
Streaming de données avec AWS Kinesis et Lambda

Assembler le tout

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"
Streaming de données avec AWS Kinesis et Lambda

Assembler le tout

res = firehose.put_record(
    DeliveryStreamName='gps-delivery-stream',
    Record = {
        'Data': payload + "\n" #<-- Saut de ligne !
    }
)
Streaming de données avec AWS Kinesis et Lambda

Fichiers créés

McYNJ.png

Streaming de données avec AWS Kinesis et Lambda

Données d'exemple

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
Streaming de données avec AWS Kinesis et Lambda

Fichiers créés

X0px9.png

Streaming de données avec AWS Kinesis et Lambda

Créer le client S3

# Créer le client S3 boto3.
s3 = boto3.client('s3', 
                  aws_access_key_id=AWS_KEY_ID, 
                  aws_secret_access_key=AWS_SECRET, 
                  region_name='us-east-1')
Streaming de données avec AWS Kinesis et Lambda

Lire les données dans un DataFrame

# Récupérer l'objet depuis S3
obj_data = s3.get_object(Bucket='sd-vehicle-data', Key=KEY_YOU_COPIED)
# Lire l'objet dans un DataFrame
vehicle_data = pd.read_csv(
    data['Body'], 
    delimiter = " ", 
    names=["record_id", "timestamp", "vin", "lon", "lat", "speed"]))
Streaming de données avec AWS Kinesis et 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
Streaming de données avec AWS Kinesis et Lambda

Récapitulatif

qUq0E.png

Streaming de données avec AWS Kinesis et Lambda

Récapitulatif

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",
}
)
Streaming de données avec AWS Kinesis et Lambda

Récapitulatif

wCJz8.png

Streaming de données avec AWS Kinesis et Lambda

Récapitulatif

Streaming de données avec AWS Kinesis et Lambda

Récapitulatif

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
Streaming de données avec AWS Kinesis et Lambda

Passons à la pratique !

Streaming de données avec AWS Kinesis et Lambda

Preparing Video For Download...