Deine erste Live-Lambda!

Datenstreaming mit AWS Kinesis und Lambda

Maksim Pecherskiy

Data Engineer

Letzte Lektion

StkGt.png

Datenstreaming mit AWS Kinesis und Lambda

Diese Lektion

Zb4Va.png

Datenstreaming mit AWS Kinesis und Lambda

Common Log Format

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
Datenstreaming mit AWS Kinesis und Lambda

S3-Dateien lesen

import pandas as pd

results in:

Response:
{
  "errorMessage": "Unable to import module 'lambda_function': No module named 'pandas'",
  "errorType": "Runtime.ImportModuleError"
}
Datenstreaming mit AWS Kinesis und Lambda

Handler aktualisieren

#lambda_function.py
import json, boto3, pandas as pd
# Clients initialisieren
...

SPEED_ALERT_THRESHOLD = 45 ALERT_PHONE_NUMBER = "+1234567890"
# Hilfsfunktion: DataFrame aus neuen Records holen def get_new_data(event): pd.read_csv()... ... return data
# Lambda-Handler def record_created_handler(event, context): data = get_new_data(event) ... sns.publish() ...
Datenstreaming mit AWS Kinesis und Lambda

get_new_data()

def get_new_data(event):
    # Liste zum Speichern neuer Objekt-Keys.
    written_objects = []

    # Über alle S3-Ereignisse iterieren.
    for record in event['Records']:

        # Zu prüfende Variablen holen
        event_name = record['eventName']
        bucket_name = record['s3']['bucket']['name']
        obj_key = record['s3']['object']['key']
Datenstreaming mit AWS Kinesis und Lambda

get_new_data()

def get_new_data(event):
        ...
        # Prüfen, ob Event aus dem Bucket sd-vehicle-data stammt.
        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)
# Neue Datensätze zu einem DataFrame zusammenführen. return pd.concat(written_objects)
Datenstreaming mit AWS Kinesis und Lambda

record_created_handler()

SPEED_ALERT_THRESHOLD = 45
ALERT_PHONE_NUMBER = "+1234567890"
...
def record_created_handler(event, context):
    # Hilfsfunktion aufrufen
    data = get_new_data(event)

## Top-Geschwindigkeiten ermitteln top_speeds = data.groupby(['vin'])['speed'].max().reset_index()
## Top-Geschwindigkeiten über 45 filtern too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD, :]
Datenstreaming mit AWS Kinesis und Lambda

record_created_handler()

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

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

## Geht noch nirgendwohin, aber wir müssen etwas zurückgeben. totals = data.groupby(['vin'])['speed'].max().reset_index() return totals.to_csv(sep=" ", index=False)
Datenstreaming mit AWS Kinesis und Lambda

Die Lambda-Funktion testen

jBpMo.png

KodWB.png

Datenstreaming mit AWS Kinesis und Lambda

Umgebungsvariablen hinzufügen

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):
    ...
Datenstreaming mit AWS Kinesis und Lambda

Umgebungsvariablen hinzufügen

7jD6x.png

Datenstreaming mit AWS Kinesis und Lambda

Umgebungsvariablen hinzufügen

eDRTX.png

Datenstreaming mit AWS Kinesis und Lambda

yKjWC.png

Datenstreaming mit AWS Kinesis und Lambda

Trigger hinzufügen

gLmbZD.png

Datenstreaming mit AWS Kinesis und Lambda

Trigger hinzufügen

7F08p.png

Datenstreaming mit AWS Kinesis und Lambda

Trigger hinzufügen

bPUtX.png

Datenstreaming mit AWS Kinesis und Lambda

Zusammenfassung

ZAbfT.png

Datenstreaming mit AWS Kinesis und Lambda

Zusammenfassung

m9jlv.png

Datenstreaming mit AWS Kinesis und Lambda

Zusammenfassung

GTlyD.png

Datenstreaming mit AWS Kinesis und Lambda

Zusammenfassung

XddOW.png

Datenstreaming mit AWS Kinesis und Lambda

Lass uns üben!

Datenstreaming mit AWS Kinesis und Lambda

Preparing Video For Download...