Une Lambda de transformation

Données en continu avec AWS Kinesis et Lambda

Maksim Pecherskiy

Data Engineer

Au chapitre 1…

2020-07-15_06-18.png

Données en continu avec AWS Kinesis et Lambda

Au chapitre 2…

2020-07-15_06-20.png

Données en continu avec AWS Kinesis et Lambda

Dans ce chapitre…

2020-07-15_06-27.png

Données en continu avec AWS Kinesis et Lambda

Approche précédente

2020-07-14_07-42.png

Données en continu avec AWS Kinesis et Lambda

Traitement dans S3 vs transformation Lambda

Traitement une fois dans S3

  • Utilise une fonction Lambda déclenchée à l'écriture d'un objet dans S3
  • Délai plus long avant la transformation
  • Données brutes stockées avant le nettoyage
  • Exige S3 comme destination Firehose

Traitement via transformation Lambda

  • Utilise une fonction Lambda déclenchée au milieu du flux Firehose
  • Transformation immédiate
  • Seules les données nettoyées sont stockées
  • Permet d'autres destinations dans Firehose
Données en continu avec AWS Kinesis et Lambda

Données entrantes

record_id timestamp vin lon lat speed
939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 4FTEX4944AK844294 106.9447146 -6.3385652 25
f29a5b3d-d0fa-43c0-9e1a-e2a5cdb8be7a 8:10:47.000 3FTEX1G5XAK844393 108.580681 34.79925 37
ff8e7131-408d-463b-8d07-d016419b0656 20:26:44.000 2LAXX1C8XAK844292 114.392392 36.097577 90
Données en continu avec AWS Kinesis et Lambda

Lambda de transformation

2020-07-15_06-38.png

Données en continu avec AWS Kinesis et Lambda

Événement d'exemple

{
  "invocationId": "invocationIdExample",
  "deliveryStreamArn": "arn:aws:firehose:us-east-1:458912630:deliverystream/gps-delivery-stream",
  "region": "us-east-1",
  "records": [
    {

"recordId": "49546986683135544286507457936321625675700192471156785154",
"approximateArrivalTimestamp": 1495072949453,
"data": "NjQuMjQyLjg4LjEwIC0gLSBbMDcvTWFyLzIwMDQ6MTY6MTA6MDIgLTA4MDBdICJHRVQgL21haWxtYW4vbGlzdGluZm8vaHNkaXZpc2lvbiBIVFRQLzEuMSIgMjAwIDYyOTE=="
} ] }
Données en continu avec AWS Kinesis et Lambda

Base64

2020-07-19_10-54.png

Données en continu avec AWS Kinesis et Lambda

Traitement des données

import base64
from datetime import datetime as dt

def convert_timestamp(record_time_val):
# Obtenir la date d'aujourd'hui en chaîne today = dt.today().strftime("%Y-%m-%d")
# Combiner la date du jour avec l'heure de l'enregistrement et créer un objet datetime new_ts = dt.strptime(f"{today} {record_time_val}", "%Y-%m-%d %H:%M:%S.%f")
# Convertir l'objet datetime en chaîne bien formatée return new_ts.strftime("%Y-%m-%dT%H:%M:%S")
Données en continu avec AWS Kinesis et Lambda

Traitement des données

def lambda_handler(event, context):
    output = []

for record in event['records']: payload = base64.b64decode(record['data'])
payload = payload.decode()
payload = payload.split(" ")
payload[1] = convert_timestamp(payload[1])
Données en continu avec AWS Kinesis et Lambda

Traitement des données

    for record in event['records']:
        ...
        payload = " ".join(payload)
        payload_enc = base64.b64encode(payload.encode())

output.append({
'recordId': record['recordId'],
'result': 'Ok',
'data': payload_enc
})
return {'records': output}
Données en continu avec AWS Kinesis et Lambda

En résumé

import base64
def lambda_handler(event, context):

output = [] for record in event['records']: # Parcourir les enregistrements
payload = base64.b64decode(record['data']).decode() # Décoder la charge utile
# La modifier
payload_enc = base64.b64encode(payload.encode()) # Réencoder la charge utile
output.append({ # L'ajouter dans un dictionnaire 'recordId': record['recordId'], 'result': 'Ok', 'data': payload_enc, })
return {'records': output}
Données en continu avec AWS Kinesis et Lambda

Créer une Lambda dans AWS

Données en continu avec AWS Kinesis et Lambda

Script pour créer la Lambda

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