Lambda trực tiếp đầu tiên của bạn!

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

Maksim Pecherskiy

Data Engineer

Bài trước

StkGt.png

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

Bài này

Zb4Va.png

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

Định dạng log phổ biến

bc75da5f-1bf6-444c-80ad-49c180e1b8de 23:16:06.000 3FTEX1G5XAK844393 -76.6990172 2.481207 40
ff8e7131-408d-463b-8d07-d016419b0656 20:26:44.000 2LAXX1C8XAK844292 114.39239199999999 36.097577 90
f29a5b3d-d0fa-43c0-9e1a-e2a5cdb8be7a 8:10:47.000 3FTEX1G5XAK844393 108.58068100000001 34.79925 37
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Đọc file S3

import pandas as pd

results in:

Response:
{
  "errorMessage": "Unable to import module 'lambda_function': No module named 'pandas'",
  "errorType": "Runtime.ImportModuleError"
}
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Cập nhật handler

#lambda_function.py
import json, boto3, pandas as pd
# Khởi tạo client
...

SPEED_ALERT_THRESHOLD = 45 ALERT_PHONE_NUMBER = "+1234567890"
# Hàm trợ giúp lấy dataframe từ bản ghi mới def get_new_data(event): pd.read_csv()... ... return data
# Trình xử lý Lambda def record_created_handler(event, context): data = get_new_data(event) ... sns.publish() ...
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

get_new_data()

def get_new_data(event):
    # Tạo danh sách lưu khóa đối tượng mới.
    written_objects = []

    # Lặp qua từng bản ghi sự kiện S3.
    for record in event['Records']:

        # Lấy các biến cần kiểm tra
        event_name = record['eventName']
        bucket_name = record['s3']['bucket']['name']
        obj_key = record['s3']['object']['key']
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

get_new_data()

def get_new_data(event):
        ...
        # Xác minh sự kiện đến từ bucket sd-vehicle-data.
        if event_name == 'ObjectCreated:Put' and bucket_name == 'sd-vehicle-data':

obj = s3.get_object(Bucket=bucket_name, Key = obj_key) df = pd.read_csv(obj['Body'], delimiter = " ", names=["record_id", "timestamp", "vin", "lon", "lat", "speed"])
written_objects.append(df)
# Gộp các bản ghi mới thành một dataframe. return pd.concat(written_objects)
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

record_created_handler()

SPEED_ALERT_THRESHOLD = 45
ALERT_PHONE_NUMBER = "+1234567890"
...
def record_created_handler(event, context):
    # Gọi hàm trợ giúp
    data = get_new_data(event)

## Lấy tốc độ cao nhất top_speeds = data.groupby(['vin'])['speed'].max().reset_index()
## Lọc tốc độ vượt ngưỡng 45 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD, :]
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

record_created_handler()

SPEED_ALERT_THRESHOLD = 45
ALERT_PHONE_NUMBER = "+1234567890"
...

def record_created_handler(event, context):
    ...
    ## Gửi SMS
    sns.publish(PhoneNumber=ALERT_PHONE_NUMBER,
        Message="Speeding Alert \n" + too_fast.to_string())

## Chưa dùng đến, nhưng cần trả về gì đó. totals = data.groupby(['vin'])['speed'].max().reset_index() return totals.to_csv(sep=" ", index=False)
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Kiểm thử hàm Lambda

jBpMo.png

KodWB.png

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

Thêm biến môi trường

import os
os.environ.get("ENV_VARIABLE_NAME", "DEFAULT_VALUE")
import os
SPEED_ALERT_THRESHOLD = os.environ.get("SPEED_ALERT_THRESHOLD", 45)
ALERT_PHONE_NUMBER = os.environ.get("ALERT_PHONE_NUMBER", None)
...

def record_created_handler(event, context):
    ...
Xử lý dữ liệu streaming với AWS Kinesis và Lambda

Thêm biến môi trường

7jD6x.png

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

Thêm biến môi trường

eDRTX.png

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

yKjWC.png

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

Thêm trigger

gLmbZD.png

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

Thêm trigger

7F08p.png

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

Thêm trigger

bPUtX.png

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

Ôn tập

ZAbfT.png

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

Ôn tập

m9jlv.png

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

Ôn tập

GTlyD.png

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

Ôn tập

XddOW.png

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...