Votre première lambda en direct !

Données en continu avec AWS Kinesis et Lambda

Maksim Pecherskiy

Data Engineer

Leçon précédente

StkGt.png

Données en continu avec AWS Kinesis et Lambda

Cette leçon

Zb4Va.png

Données en continu avec AWS Kinesis et Lambda

Format de journal courant

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

Lecture de fichiers S3

import pandas as pd

résulte en :

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

Mise à jour du gestionnaire

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

SPEED_ALERT_THRESHOLD = 45 ALERT_PHONE_NUMBER = "+1234567890"
# Helper function to get dataframe from written records def get_new_data(event): pd.read_csv()... ... return data
# Lambda function handler def record_created_handler(event, context): data = get_new_data(event) ... sns.publish() ...
Données en continu avec AWS Kinesis et Lambda

get_new_data()

def get_new_data(event):
    # Create a list to store new object keys.
    written_objects = []

    # Iterate over each S3 event record.
    for record in event['Records']:

        # Get the variables to check for
        event_name = record['eventName']
        bucket_name = record['s3']['bucket']['name']
        obj_key = record['s3']['object']['key']
Données en continu avec AWS Kinesis et Lambda

get_new_data()

def get_new_data(event):
        ...
        # Verify that event is created from sd-vehicle-data bucket.
        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)
# Concatenate new records into a single dataframe. return pd.concat(written_objects)
Données en continu avec AWS Kinesis et Lambda

record_created_handler()

SPEED_ALERT_THRESHOLD = 45
ALERT_PHONE_NUMBER = "+1234567890"
...
def record_created_handler(event, context):
    # Call the helper method
    data = get_new_data(event)

## Get the top speeds top_speeds = data.groupby(['vin'])['speed'].max().reset_index()
## Get top speeds that exceed the limit of 45 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD, :]
Données en continu avec AWS Kinesis et Lambda

record_created_handler()

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

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

## This doesn't go anywhere yet, but we need to return something. totals = data.groupby(['vin'])['speed'].max().reset_index() return totals.to_csv(sep=" ", index=False)
Données en continu avec AWS Kinesis et Lambda

Tester la fonction Lambda

jBpMo.png

KodWB.png

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

Ajout de variables d'environnement

7jD6x.png

Données en continu avec AWS Kinesis et Lambda

Ajout de variables d'environnement

eDRTX.png

Données en continu avec AWS Kinesis et Lambda

yKjWC.png

Données en continu avec AWS Kinesis et Lambda

Ajout d'un déclencheur

gLmbZD.png

Données en continu avec AWS Kinesis et Lambda

Ajout d'un déclencheur

7F08p.png

Données en continu avec AWS Kinesis et Lambda

Ajout d'un déclencheur

bPUtX.png

Données en continu avec AWS Kinesis et Lambda

Révision

ZAbfT.png

Données en continu avec AWS Kinesis et Lambda

Révision

m9jlv.png

Données en continu avec AWS Kinesis et Lambda

Révision

GTlyD.png

Données en continu avec AWS Kinesis et Lambda

Révision

XddOW.png

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