Quy trình dữ liệu serverless

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Maksim Pecherskiy

Data Engineer

Ôn lại

z8HRk.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Phân vùng hiện tại

FMktI.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Lọc người vượt tốc

IPwGK.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tổng hợp theo ngày

tvBXx.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Mở recordReaderS3

aIahM.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Chỉnh sửa recordReaderS3

...
import pytz
from datetime import datetime
tz = pytz.timezone('America/Los_Angeles')
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Chỉnh sửa recordReaderS3

def record_created_handler(event, context):

... ## Lấy các tốc độ vượt quá giới hạn 45 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD :]
## Tạo khóa đối tượng fdate = datetime.now(tz).strftime("%Y%m%d/%HH%MM%SS") obj_key = f"speeders/{fdate}.csv" # tên tệp trong thư mục speeders
## Ghi đối tượng lên S3 s3.put_object(Bucket='sd-vehicle-data', Key=obj_key, Body=too_fast.to_csv(sep=" ", index=False) )
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ghi người vượt tốc theo ngày, giờ

s3QY4.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

speederAggregator

tvBXx.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tạo speederAggregator

3pGYX.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Thêm layer AWS Data Wrangler

7TO0V.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tài nguyên speederAggregator

sO0op.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tạo trigger theo thời gian

4Qi4i.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

y0FkQ.png

cron(50 6 * * ? *)
  • Vào phút 50
  • Của giờ thứ 6 (UTC)
  • Mỗi ngày trong tháng
  • Mỗi tháng
  • Mỗi ngày trong tuần
  • Mỗi năm
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Hàm gọi lại 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")
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Hàm gọi lại 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 = " ")) # Gộp bản ghi mới thành một dataframe. data = pd.concat(day_data) data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Gói 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()
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Cách cũ 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 = " "))

  # Gộp bản ghi mới thành một dataframe.
  data = pd.concat(day_data)
  data.columns = ["record_id", "timestamp", "vin", "lon", "lat", "speed"]
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Cách cũ 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()
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ghi tệp tổng hợp người vượt tốc

wr.s3.to_csv(df, 
             f"s3://sd-vehicle-data/speeders_daily/{filter_date}.csv", 
             boto3_session=session, 
             sep=" ", index=False)
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Lưu và kiểm thử

agF2b.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Bật trigger

COfYY.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tổng kết

u1Jwm.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tổng kết

7ZYPV.png

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Tổng kết

  • Cron để lập lịch
  • wr.s3.read_csv()
  • wr.s3.write_csv()
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Ayo berlatih!

Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Preparing Video For Download...