Flux de données serverless

Streaming de données avec AWS Kinesis et Lambda

Maksim Pecherskiy

Data Engineer

Retour en arrière

z8HRk.png

Streaming de données avec AWS Kinesis et Lambda

Partitionnement actuel

FMktI.png

Streaming de données avec AWS Kinesis et Lambda

Filtrer les contrevenants

IPwGK.png

Streaming de données avec AWS Kinesis et Lambda

Agrégation par jour

tvBXx.png

Streaming de données avec AWS Kinesis et Lambda

Ouvrir recordReaderS3

aIahM.png

Streaming de données avec AWS Kinesis et Lambda

Modification de recordReaderS3

...
import pytz
from datetime import datetime
tz = pytz.timezone('America/Los_Angeles')
Streaming de données avec AWS Kinesis et Lambda

Modification de recordReaderS3

def record_created_handler(event, context):

... ## Obtenir les vitesses au-delà de 45 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD :]
## Générer la clé d'objet fdate = datetime.now(tz).strftime("%Y%m%d/%HH%MM%SS") obj_key = f"speeders/{fdate}.csv" # nom de fichier dans le dossier speeders
## Écrire l'objet dans S3 s3.put_object(Bucket='sd-vehicle-data', Key=obj_key, Body=too_fast.to_csv(sep=" ", index=False) )
Streaming de données avec AWS Kinesis et Lambda

Écriture des contrevenants par date et heure

s3QY4.png

Streaming de données avec AWS Kinesis et Lambda

speederAggregator

tvBXx.png

Streaming de données avec AWS Kinesis et Lambda

Créer speederAggregator

3pGYX.png

Streaming de données avec AWS Kinesis et Lambda

Ajouter la couche AWS Data Wrangler

7TO0V.png

Streaming de données avec AWS Kinesis et Lambda

Ressources de speederAggregator

sO0op.png

Streaming de données avec AWS Kinesis et Lambda

Créer le déclencheur programmé

4Qi4i.png

Streaming de données avec AWS Kinesis et Lambda

y0FkQ.png

cron(50 6 * * ? *)
  • À la minute 50
  • À la 6e heure (UTC)
  • Chaque jour du mois
  • Chaque mois
  • Chaque jour de la semaine
  • Chaque année
Streaming de données avec AWS Kinesis et Lambda

Callback 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")
Streaming de données avec AWS Kinesis et Lambda

Callback 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"]
Streaming de données avec AWS Kinesis et Lambda

Package 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()
Streaming de données 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"]
Streaming de données 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()
Streaming de données avec AWS Kinesis et Lambda

Écriture du fichier agrégé des contrevenants

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

Enregistrer et tester

agF2b.png

Streaming de données avec AWS Kinesis et Lambda

Activer le déclencheur

COfYY.png

Streaming de données avec AWS Kinesis et Lambda

Récapitulatif

u1Jwm.png

Streaming de données avec AWS Kinesis et Lambda

Récapitulatif

7ZYPV.png

Streaming de données avec AWS Kinesis et Lambda

Récapitulatif

  • Cron pour la planification
  • wr.s3.read_csv()
  • wr.s3.write_csv()
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...