Votre première Lambda en direct !

Streaming de données avec AWS Kinesis et Lambda

Maksim Pecherskiy

Data Engineer

Leçon précédente

Récap précédent

Streaming de données avec AWS Kinesis et Lambda

Cette leçon

Plan de la leçon

Streaming de données avec AWS Kinesis et Lambda

Format de journal commun

bc75da5f-1bf6-444c-80ad-49c180e1b8de 23:16:06.000 3FTEX1G5XAK844393 -76.6990172 2.481207 40
ff8e7131-408d-463b-8d07-d016419b0656 20:26:44.000 2LAXX1C8XAK844292 114.39239199999999 36.097577 90
f29a5b3d-d0fa-43c0-9e1a-e2a5cdb8be7a 8:10:47.000 3FTEX1G5XAK844393 108.58068100000001 34.79925 37
Streaming de données avec AWS Kinesis et Lambda

Lecture de fichiers S3

import pandas as pd

results in:

Response:
{
  "errorMessage": "Unable to import module 'lambda_function': No module named 'pandas'",
  "errorType": "Runtime.ImportModuleError"
}
Streaming de données avec AWS Kinesis et Lambda

Mise à jour du handler

#lambda_function.py
import json, boto3, pandas as pd
# Initialiser les clients
...

SPEED_ALERT_THRESHOLD = 45 ALERT_PHONE_NUMBER = "+1234567890"
# Fonction utilitaire pour lire un dataframe depuis les enregistrements écrits def get_new_data(event): pd.read_csv()... ... return data
# Gestionnaire Lambda def record_created_handler(event, context): data = get_new_data(event) ... sns.publish() ...
Streaming de données avec AWS Kinesis et Lambda

get_new_data()

def get_new_data(event):
    # Créer une liste pour stocker les clés des nouveaux objets.
    written_objects = []

    # Parcourir chaque enregistrement d'événement S3.
    for record in event['Records']:

        # Récupérer les variables à vérifier
        event_name = record['eventName']
        bucket_name = record['s3']['bucket']['name']
        obj_key = record['s3']['object']['key']
Streaming de données avec AWS Kinesis et Lambda

get_new_data()

def get_new_data(event):
        ...
        # Vérifier que l'événement vient du bucket sd-vehicle-data.
        if event_name == 'ObjectCreated:Put' and bucket_name == 'sd-vehicle-data':

obj = s3.get_object(Bucket=bucket_name, Key = obj_key) df = pd.read_csv(obj['Body'], delimiter = " ", names=["record_id", "timestamp", "vin", "lon", "lat", "speed"])
written_objects.append(df)
# Concaténer les nouveaux enregistrements dans un seul dataframe. return pd.concat(written_objects)
Streaming de données avec AWS Kinesis et Lambda

record_created_handler()

SPEED_ALERT_THRESHOLD = 45
ALERT_PHONE_NUMBER = "+1234567890"
...
def record_created_handler(event, context):
    # Appeler la fonction utilitaire
    data = get_new_data(event)

## Obtenir les vitesses max top_speeds = data.groupby(['vin'])['speed'].max().reset_index()
## Filtrer celles > 45 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD, :]
Streaming de données avec AWS Kinesis et Lambda

record_created_handler()

SPEED_ALERT_THRESHOLD = 45
ALERT_PHONE_NUMBER = "+1234567890"
...

def record_created_handler(event, context):
    ...
    ## Envoyer un SMS
    sns.publish(PhoneNumber=ALERT_PHONE_NUMBER,
        Message="Speeding Alert \n" + too_fast.to_string())

## Rien n'est encore renvoyé, mais il faut retourner quelque chose. totals = data.groupby(['vin'])['speed'].max().reset_index() return totals.to_csv(sep=" ", index=False)
Streaming de données avec AWS Kinesis et Lambda

Tester la fonction Lambda

Sortie de test 1

Sortie de test 2

Streaming de données avec AWS Kinesis et Lambda

Ajout de variables d'environnement

import os
os.environ.get("ENV_VARIABLE_NAME", "DEFAULT_VALUE")
import os
SPEED_ALERT_THRESHOLD = os.environ.get("SPEED_ALERT_THRESHOLD", 45)
ALERT_PHONE_NUMBER = os.environ.get("ALERT_PHONE_NUMBER", None)
...

def record_created_handler(event, context):
    ...
Streaming de données avec AWS Kinesis et Lambda

Ajout de variables d'environnement

Définir variables

Streaming de données avec AWS Kinesis et Lambda

Ajout de variables d'environnement

Configurer variables

Streaming de données avec AWS Kinesis et Lambda

Ajouter un déclencheur

Streaming de données avec AWS Kinesis et Lambda

Ajout d'un déclencheur

Sélection du déclencheur

Streaming de données avec AWS Kinesis et Lambda

Ajout d'un déclencheur

Configuration du déclencheur

Streaming de données avec AWS Kinesis et Lambda

Ajout d'un déclencheur

Règles S3

Streaming de données avec AWS Kinesis et Lambda

Bilan

Résumé 1

Streaming de données avec AWS Kinesis et Lambda

Bilan

Résumé 2

Streaming de données avec AWS Kinesis et Lambda

Bilan

Résumé 3

Streaming de données avec AWS Kinesis et Lambda

Bilan

Résumé 4

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...