एक परिवर्तनकारी Lambda

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

Maksim Pecherskiy

Data Engineer

अध्याय 1 में...

2020-07-15_06-18.png

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

अध्याय 2 में...

2020-07-15_06-20.png

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

इस अध्याय में...

2020-07-15_06-27.png

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

पिछला तरीका

2020-07-14_07-42.png

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

S3 में एक बार बनाम Lambda ट्रांसफॉर्म

S3 में एक बार प्रोसेसिंग

  • S3 में ऑब्जेक्ट लिखने पर ट्रिगर होने वाला lambda फंक्शन उपयोग करता है
  • डेटा ट्रांसफॉर्म होने तक अधिक देरी
  • क्लीनिंग से पहले कच्चा डेटा स्टोर होता है
  • Firehose का डेस्टिनेशन S3 होना जरूरी

Lambda ट्रांसफॉर्म के जरिये प्रोसेसिंग

  • Firehose स्ट्रीम के बीच में ट्रिगर होने वाला lambda फंक्शन उपयोग करता है
  • ट्रांसफॉर्मेशन तुरंत होता है
  • केवल साफ किया हुआ डेटा स्टोर होता है
  • Firehose में दूसरे डेस्टिनेशन भी संभव
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

आने वाला डेटा

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
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

Transformational Lambda

2020-07-15_06-38.png

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

नमूना इवेंट

{
  "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=="
} ] }
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

Base64

2020-07-19_10-54.png

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

डेटा प्रोसेस करना

import base64
from datetime import datetime as dt

def convert_timestamp(record_time_val):
# आज की तारीख स्ट्रिंग रूप में लें today = dt.today().strftime("%Y-%m-%d")
# आज की तारीख को रिकॉर्ड के समय से जोड़कर datetime ऑब्जेक्ट बनाएँ new_ts = dt.strptime(f"{today} {record_time_val}", "%Y-%m-%d %H:%M:%S.%f")
# datetime ऑब्जेक्ट को अच्छे फॉर्मेट की स्ट्रिंग में बदलें return new_ts.strftime("%Y-%m-%dT%H:%M:%S")
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

डेटा प्रोसेस करना

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])
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

डेटा प्रोसेस करना

    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}
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

एक रिव्यू

import base64
def lambda_handler(event, context):

output = [] for record in event['records']: # रिकॉर्ड्स पर इटरेट करें
payload = base64.b64decode(record['data']).decode() # पेलोड डिकोड करें
# इसे संशोधित करें
payload_enc = base64.b64encode(payload.encode()) # दोबारा एन्कोड करें
output.append({ # इसे डिक्शनरी में डालें 'recordId': record['recordId'], 'result': 'Ok', 'data': payload_enc, })
return {'records': output}
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

AWS में Lambda बनाना

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

Lambda बनाने की स्क्रिप्ट

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

अभ्यास करते हैं!

AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

Preparing Video For Download...