Serverloser Daten-Workflow

Datenstreaming mit AWS Kinesis und Lambda

Maksim Pecherskiy

Data Engineer

Rückblick

z8HRk.png

Datenstreaming mit AWS Kinesis und Lambda

Aktuelle Partitionierung

FMktI.png

Datenstreaming mit AWS Kinesis und Lambda

Speeders herausfiltern

IPwGK.png

Datenstreaming mit AWS Kinesis und Lambda

Täglich aggregieren

tvBXx.png

Datenstreaming mit AWS Kinesis und Lambda

recordReaderS3 öffnen

aIahM.png

Datenstreaming mit AWS Kinesis und Lambda

recordReaderS3 bearbeiten

...
import pytz
from datetime import datetime
tz = pytz.timezone('America/Los_Angeles')
Datenstreaming mit AWS Kinesis und Lambda

recordReaderS3 bearbeiten

def record_created_handler(event, context):

... ## Top-Geschwindigkeiten über dem Limit 45 abrufen too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD :]
## Objekt-Schlüssel erzeugen fdate = datetime.now(tz).strftime("%Y%m%d/%HH%MM%SS") obj_key = f"speeders/{fdate}.csv" # Dateiname im Ordner speeders
## Objekt nach S3 schreiben s3.put_object(Bucket='sd-vehicle-data', Key=obj_key, Body=too_fast.to_csv(sep=" ", index=False) )
Datenstreaming mit AWS Kinesis und Lambda

Speeders nach Datum und Uhrzeit schreiben

s3QY4.png

Datenstreaming mit AWS Kinesis und Lambda

speederAggregator

tvBXx.png

Datenstreaming mit AWS Kinesis und Lambda

speederAggregator erstellen

3pGYX.png

Datenstreaming mit AWS Kinesis und Lambda

AWS Data Wrangler-Layer hinzufügen

7TO0V.png

Datenstreaming mit AWS Kinesis und Lambda

Ressourcen für speederAggregator

sO0op.png

Datenstreaming mit AWS Kinesis und Lambda

Zeitgesteuerten Trigger erstellen

4Qi4i.png

Datenstreaming mit AWS Kinesis und Lambda

y0FkQ.png

cron(50 6 * * ? *)
  • In Minute 50
  • Der 6. Stunde (UTC)
  • An jedem Tag des Monats
  • In jedem Monat
  • An jedem Wochentag
  • In jedem Jahr
Datenstreaming mit AWS Kinesis und Lambda

speederAggregator-Callback

import boto3, pytz, pandas as pd

s3 = boto3.client("s3" ...)
...

def speeder_aggregator(event, context): tz = pytz.timezone('America/Los_Angeles') filter_date = datetime.now(tz).strftime("%Y%m%d")
Datenstreaming mit AWS Kinesis und Lambda

speederAggregator-Callback

def speeder_aggregator(event, context):
  ...
  objects = s3.list_objects_v2(
      Bucket='sd-vehicle-data', Prefix=f'speeders/{filter_date}')
  day_data = []

for obj in objects['Contents']: print(obj['Key']) day_record = s3.get_object(Bucket='sd-vehicle-data', Key = obj['Key']) day_data.append( pd.read_csv(day_record['Body'], delimiter = " ")) # Neue Datensätze zu einem DataFrame zusammenführen. data = pd.concat(day_data) data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
Datenstreaming mit AWS Kinesis und Lambda

awswrangler-Paket

import awswrangler as wr

session = boto3.Session(aws_access_key_id = AWS_KEY, aws_secret_access_key = AWS_SECRET, region_name="us-east-1")
df = wr.s3.read_csv(f"s3://sd-vehicle-data/speeders/{filter_date}") df.head()
Datenstreaming mit AWS Kinesis und Lambda

Alt vs. awswrangler

def speeder_aggregator(event, context):
  ...
  objects = s3.list_objects_v2(
      Bucket='sd-vehicle-data', Prefix=f'speeders/{filter_date}')
  day_data = []

  for obj in objects['Contents']:
      print(obj['Key'])
      day_record = s3.get_object(Bucket='sd-vehicle-data', Key = obj['Key'])
      day_data.append(
          pd.read_csv(day_record['Body'], delimiter = " "))

  # Neue Datensätze zu einem DataFrame zusammenführen.
  data = pd.concat(day_data)
  data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
Datenstreaming mit AWS Kinesis und Lambda

Alt vs. awswrangler

import awswrangler as wr

session = boto3.Session(aws_access_key_id = AWS_KEY, 
                        aws_secret_access_key = AWS_SECRET, region_name="us-east-1")
def speeder_aggregator(event, context):
    df = wr.s3.read_csv(f"s3://sd-vehicle-data/speeders/{filter_date}")      
    df.head()
Datenstreaming mit AWS Kinesis und Lambda

Aggregierte Speeders-Datei schreiben

wr.s3.to_csv(df, 
             f"s3://sd-vehicle-data/speeders_daily/{filter_date}.csv", 
             boto3_session=session, 
             sep=" ", index=False)
Datenstreaming mit AWS Kinesis und Lambda

Speichern und testen

agF2b.png

Datenstreaming mit AWS Kinesis und Lambda

Trigger aktivieren

COfYY.png

Datenstreaming mit AWS Kinesis und Lambda

Review

u1Jwm.png

Datenstreaming mit AWS Kinesis und Lambda

Review

7ZYPV.png

Datenstreaming mit AWS Kinesis und Lambda

Review

  • Cron zum Planen
  • wr.s3.read_csv()
  • wr.s3.write_csv()
Datenstreaming mit AWS Kinesis und Lambda

Lass uns üben!

Datenstreaming mit AWS Kinesis und Lambda

Preparing Video For Download...