Flux de données serverless

Données en continu avec AWS Kinesis et Lambda

Maksim Pecherskiy

Data Engineer

Retour en arrière

z8HRk.png

Données en continu avec AWS Kinesis et Lambda

Partitionnement actuel

FMktI.png

Données en continu avec AWS Kinesis et Lambda

Filtrer les chauffards

IPwGK.png

Données en continu avec AWS Kinesis et Lambda

Agrégation par jour

tvBXx.png

Données en continu avec AWS Kinesis et Lambda

Ouvrir recordReaderS3

aIahM.png

Données en continu avec AWS Kinesis et Lambda

Modifier recordReaderS3

...
import pytz
from datetime import datetime
tz = pytz.timezone('America/Los_Angeles')
Données en continu avec AWS Kinesis et Lambda

Modifier recordReaderS3

def record_created_handler(event, context):

... ## Get top speeds that exceed the limit of 45 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD :]
## Generate object key fdate = datetime.now(tz).strftime("%Y%m%d/%HH%MM%SS") obj_key = f"speeders/{fdate}.csv" # filename in speeders folder
## Write the object to S3 s3.put_object(Bucket='sd-vehicle-data', Key=obj_key, Body=too_fast.to_csv(sep=" ", index=False) )
Données en continu avec AWS Kinesis et Lambda

Écrire les chauffards par date et heure

s3QY4.png

Données en continu avec AWS Kinesis et Lambda

speederAggregator

tvBXx.png

Données en continu avec AWS Kinesis et Lambda

Créer speederAggregator

3pGYX.png

Données en continu avec AWS Kinesis et Lambda

Ajouter la couche AWS Data Wrangler

7TO0V.png

Données en continu avec AWS Kinesis et Lambda

Ressources de speederAggregator

sO0op.png

Données en continu avec AWS Kinesis et Lambda

Créer le déclencheur temporisé

4Qi4i.png

Données en continu avec AWS Kinesis et Lambda

y0FkQ.png

cron(50 6 * * ? *)
  • À la minute 50
  • De la 6e heure (UTC)
  • De chaque jour du mois
  • De chaque mois
  • De chaque jour de la semaine
  • De chaque année
Données en continu avec AWS Kinesis et Lambda

Rappel de speederAggregator

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

Rappel de speederAggregator

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 = " ")) # Concatenate new records into a single dataframe. data = pd.concat(day_data) data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
Données en continu avec AWS Kinesis et Lambda

Forfait 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")
df = wr.s3.read_csv(f"s3://sd-vehicle-data/speeders/{filter_date}") df.head()
Données en continu avec AWS Kinesis et Lambda

Ancienne méthode 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 = " "))

  # Concatenate new records into a single dataframe.
  data = pd.concat(day_data)
  data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
Données en continu avec AWS Kinesis et Lambda

Ancienne méthode 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()
Données en continu avec AWS Kinesis et Lambda

Écrire le fichier agrégé des chauffards

wr.s3.to_csv(df, 
             f"s3://sd-vehicle-data/speeders_daily/{filter_date}.csv", 
             boto3_session=session, 
             sep=" ", index=False)
Données en continu avec AWS Kinesis et Lambda

Enregistrer et tester

agF2b.png

Données en continu avec AWS Kinesis et Lambda

Activer le déclencheur

COfYY.png

Données en continu avec AWS Kinesis et Lambda

Révision

u1Jwm.png

Données en continu avec AWS Kinesis et Lambda

Révision

7ZYPV.png

Données en continu avec AWS Kinesis et Lambda

Révision

  • Cron pour l'ordonnancement
  • wr.s3.read_csv()
  • wr.s3.write_csv()
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...