सर्वरलेस डेटा वर्कफ़्लो

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

Maksim Pecherskiy

Data Engineer

एक नज़र पीछे

z8HRk.png

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

वर्तमान पार्टिशनिंग

FMktI.png

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

स्पीडर फ़िल्टर करना

IPwGK.png

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

दिन के अनुसार एग्रीगेट करना

tvBXx.png

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

recordReaderS3 खोलें

aIahM.png

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

recordReaderS3 संपादित करना

...
import pytz
from datetime import datetime
tz = pytz.timezone('America/Los_Angeles')
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

recordReaderS3 संपादित करना

def record_created_handler(event, context):

... ## Get top speeds that exceed the limit of 45 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD :]
## Generate object key fdate = datetime.now(tz).strftime("%Y%m%d/%HH%MM%SS") obj_key = f"speeders/{fdate}.csv" # filename in speeders folder
## Write the object to S3 s3.put_object(Bucket='sd-vehicle-data', Key=obj_key, Body=too_fast.to_csv(sep=" ", index=False) )
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

तारीख-समय के अनुसार स्पीडर लिखना

s3QY4.png

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

speederAggregator

tvBXx.png

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

speederAggregator बनाएँ

3pGYX.png

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

AWS Data Wrangler लेयर जोड़ें

7TO0V.png

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

speederAggregator संसाधन

sO0op.png

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

टाइम्ड ट्रिगर बनाएँ

4Qi4i.png

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

y0FkQ.png

cron(50 6 * * ? *)
  • Minute 50 पर
  • 6वें घंटे (UTC) पर
  • हर महीने के हर दिन
  • हर महीने
  • हफ्ते के हर दिन
  • हर साल
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

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

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

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

पुराना तरीका बनाम 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"]
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

पुराना तरीका बनाम 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()
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

एग्रीगेटेड स्पीडर फ़ाइल लिखना

wr.s3.to_csv(df, 
             f"s3://sd-vehicle-data/speeders_daily/{filter_date}.csv", 
             boto3_session=session, 
             sep=" ", index=False)
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

सेव करें और टेस्ट करें

agF2b.png

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

ट्रिगर सक्षम करें

COfYY.png

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

समीक्षा

u1Jwm.png

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

समीक्षा

7ZYPV.png

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

समीक्षा

  • शेड्यूलिंग के लिए Cron
  • wr.s3.read_csv()
  • wr.s3.write_csv()
AWS Kinesis और Lambda के साथ स्ट्रीमिंग डेटा

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

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

Preparing Video For Download...