Serverlöst dataarbetsflöde

Strömdata med AWS Kinesis och Lambda

Maksim Pecherskiy

Data Engineer

Tillbakablick

z8HRk.png

Strömdata med AWS Kinesis och Lambda

Nuvarande partitionering

FMktI.png

Strömdata med AWS Kinesis och Lambda

Filtrera ut fortkörare

IPwGK.png

Strömdata med AWS Kinesis och Lambda

Aggregera per dag

tvBXx.png

Strömdata med AWS Kinesis och Lambda

Öppna recordReaderS3

aIahM.png

Strömdata med AWS Kinesis och Lambda

Redigera recordReaderS3

...
import pytz
from datetime import datetime
tz = pytz.timezone('America/Los_Angeles')
Strömdata med AWS Kinesis och Lambda

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

Skriva fortkörare per datum och tid

s3QY4.png

Strömdata med AWS Kinesis och Lambda

speederAggregator

tvBXx.png

Strömdata med AWS Kinesis och Lambda

Skapa speederAggregator

3pGYX.png

Strömdata med AWS Kinesis och Lambda

Lägg till AWS data wrangler-lager

7TO0V.png

Strömdata med AWS Kinesis och Lambda

speederAggregator-resurser

sO0op.png

Strömdata med AWS Kinesis och Lambda

Skapa den tidsstyrda utlösaren

4Qi4i.png

Strömdata med AWS Kinesis och Lambda

y0FkQ.png

cron(50 6 * * ? *)
  • Vid minut 50
  • Den sjätte timmen (UTC)
  • Varje dag i månaden
  • Varje månad
  • Varje veckodag
  • Varje år
Strömdata med AWS Kinesis och Lambda

speederAggregator-callback

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

speederAggregator-callback

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

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

Gamla metoden vs 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"]
Strömdata med AWS Kinesis och Lambda

Gamla metoden vs 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()
Strömdata med AWS Kinesis och Lambda

Skriva aggregerad fortkörningsfil

wr.s3.to_csv(df, 
             f"s3://sd-vehicle-data/speeders_daily/{filter_date}.csv", 
             boto3_session=session, 
             sep=" ", index=False)
Strömdata med AWS Kinesis och Lambda

Spara och testa

agF2b.png

Strömdata med AWS Kinesis och Lambda

Aktivera utlösaren

COfYY.png

Strömdata med AWS Kinesis och Lambda

Genomgång

u1Jwm.png

Strömdata med AWS Kinesis och Lambda

Genomgång

7ZYPV.png

Strömdata med AWS Kinesis och Lambda

Genomgång

  • Cron för schemaläggning
  • wr.s3.read_csv()
  • wr.s3.write_csv()
Strömdata med AWS Kinesis och Lambda

Nu kör vi en övning!

Strömdata med AWS Kinesis och Lambda

Preparing Video For Download...