Travailler avec le flux de diffusion Firehose

Données en continu avec AWS Kinesis et Lambda

Maksim Pecherskiy

Data Engineer

Prêt à créer le flux

KF414O.png

Données en continu avec AWS Kinesis et Lambda

Prêt à créer le flux

YgVbe.png

Données en continu avec AWS Kinesis et Lambda

Prêt à créer le flux

BBYB4.png

Données en continu avec AWS Kinesis et Lambda

Prêt à créer le flux

Ue73e.png

Données en continu avec AWS Kinesis et Lambda

Prêt à créer le flux

qUq0E.png

Données en continu avec AWS Kinesis et Lambda

Obtenir le rôle ARN

qH99c.png

Données en continu 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')
Données en continu 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"
}
)
Données en continu avec AWS Kinesis et Lambda

Réponse à la création du flux

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

Flux prêt

z0Fp4.png

Données en continu avec AWS Kinesis et Lambda

Écriture dans le flux

La0zZ.png

Données en continu avec AWS Kinesis et Lambda

Matériel télématique

OGHUHC.png

Données en continu avec AWS Kinesis et Lambda

Envoi des données télématiques

wCJz8.png

Données en continu avec AWS Kinesis et Lambda

Enregistrement unique

{

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

Arrivée des enregistrements

vRg6A.png

Données en continu avec AWS Kinesis et Lambda

Autre cas d'usage

gu6oT.png

Données en continu avec AWS Kinesis et Lambda

Autre cas d'usage

XFD7E.png

Données en continu avec AWS Kinesis et Lambda

Modèles récurrents

h7hUU.png

Données en continu avec AWS Kinesis et Lambda

Envoi d'un enregistrement

res = firehose.put_record(

DeliveryStreamName='gps-delivery-stream',
Record = { 'Data': payload }
)
Données en continu avec AWS Kinesis et Lambda

Envoi d'un enregistrement

    Record = {
        'Data': payload 
    }
Données en continu 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"
Données en continu 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"
Données en continu 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"
Données en continu avec AWS Kinesis et Lambda

Assembler le tout

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

Fichiers créés

McYNJ.png

Données en continu 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
Données en continu avec AWS Kinesis et Lambda

Fichiers créés

X0px9.png

Données en continu avec AWS Kinesis et Lambda

Créer un 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')
Données en continu 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"]))
Données en continu 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
Données en continu avec AWS Kinesis et Lambda

Révision

qUq0E.png

Données en continu avec AWS Kinesis et Lambda

Révision

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",
}
)
Données en continu avec AWS Kinesis et Lambda

Révision

wCJz8.png

Données en continu avec AWS Kinesis et Lambda

Révision

Données en continu avec AWS Kinesis et Lambda

Révision

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
Données en continu avec AWS Kinesis et Lambda

Passons à la pratique !

Données en continu avec AWS Kinesis et Lambda

Preparing Video For Download...