Une Lambda de transformation

Streaming de données avec AWS Kinesis et Lambda

Maksim Pecherskiy

Data Engineer

Au chapitre 1…

2020-07-15_06-18.png

Streaming de données avec AWS Kinesis et Lambda

Au chapitre 2…

2020-07-15_06-20.png

Streaming de données avec AWS Kinesis et Lambda

Dans ce chapitre…

2020-07-15_06-27.png

Streaming de données avec AWS Kinesis et Lambda

Approche précédente

2020-07-14_07-42.png

Streaming de données 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 dans S3
  • Délai plus long avant transformation
  • Données brutes stockées avant nettoyage
  • Nécessite S3 comme destination Firehose

Traitement via transformation Lambda

  • Utilise une fonction Lambda au milieu du flux Firehose
  • Transformation immédiate
  • Seules les données nettoyées sont stockées
  • Permet d'autres destinations dans Firehose
Streaming de données 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
Streaming de données avec AWS Kinesis et Lambda

Lambda de transformation

2020-07-15_06-38.png

Streaming de données avec AWS Kinesis et Lambda

Événement 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=="
} ] }
Streaming de données avec AWS Kinesis et Lambda

Base64

2020-07-19_10-54.png

Streaming de données 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 du jour en chaîne today = dt.today().strftime("%Y-%m-%d")
# Combiner la date du jour avec l'heure du record et créer un datetime new_ts = dt.strptime(f"{today} {record_time_val}", "%Y-%m-%d %H:%M:%S.%f")
# Formater l'objet datetime en chaîne lisible return new_ts.strftime("%Y-%m-%dT%H:%M:%S")
Streaming de données 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])
Streaming de données 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}
Streaming de données avec AWS Kinesis et Lambda

Récapitulatif

import base64
def lambda_handler(event, context):

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

Créer une Lambda dans AWS

Streaming de données avec AWS Kinesis et Lambda

Script pour créer la Lambda

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