En transformerande Lambda

Strömdata med AWS Kinesis och Lambda

Maksim Pecherskiy

Data Engineer

I kapitel 1...

2020-07-15_06-18.png

Strömdata med AWS Kinesis och Lambda

I kapitel 2...

2020-07-15_06-20.png

Strömdata med AWS Kinesis och Lambda

I det här kapitlet...

2020-07-15_06-27.png

Strömdata med AWS Kinesis och Lambda

Tidigare tillvägagångssätt

2020-07-14_07-42.png

Strömdata med AWS Kinesis och Lambda

Bearbetning i S3 vs. Lambda-transformation

Bearbetning en gång i S3

  • Använder en lambdafunktion som utlöses vid objektskrivning i S3
  • Längre fördröjning innan data kan transformeras
  • Rådata lagras innan rensning
  • Kräver att Firehose-destinationen är S3

Bearbetning via Lambda-transformation

  • Använder en lambdafunktion som utlöses mitt i Firehose-strömmen
  • Transformationen sker omedelbart
  • Endast rensad data lagras
  • Tillåter andra destinationer i Firehose
Strömdata med AWS Kinesis och Lambda

Inkommande data

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
Strömdata med AWS Kinesis och Lambda

Transformerande Lambda

2020-07-15_06-38.png

Strömdata med AWS Kinesis och Lambda

Exempelhändelse

{
  "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=="
} ] }
Strömdata med AWS Kinesis och Lambda

Base64

2020-07-19_10-54.png

Strömdata med AWS Kinesis och Lambda

Bearbetning av data

import base64
from datetime import datetime as dt

def convert_timestamp(record_time_val):
# Get today's date as string today = dt.today().strftime("%Y-%m-%d")
# Combine today's date with the record's time and make datetime object new_ts = dt.strptime(f"{today} {record_time_val}", "%Y-%m-%d %H:%M:%S.%f")
# Convert the datetime object to a nicely formatted string return new_ts.strftime("%Y-%m-%dT%H:%M:%S")
Strömdata med AWS Kinesis och Lambda

Bearbetning av data

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])
Strömdata med AWS Kinesis och Lambda

Bearbetning av data

    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}
Strömdata med AWS Kinesis och Lambda

En genomgång

import base64
def lambda_handler(event, context):

output = [] for record in event['records']: # Iterate over the records
payload = base64.b64decode(record['data']).decode() # Decode the payload
# Modify it
payload_enc = base64.b64encode(payload.encode()) # Re-encode it
output.append({ # Put it in a dictionary 'recordId': record['recordId'], 'result': 'Ok', 'data': payload_enc, })
return {'records': output}
Strömdata med AWS Kinesis och Lambda

Skapa Lambda i AWS

Strömdata med AWS Kinesis och Lambda

Skript för att skapa Lambda

Strömdata med AWS Kinesis och Lambda

Nu kör vi en övning!

Strömdata med AWS Kinesis och Lambda

Preparing Video For Download...