你的第一个实时 Lambda!

使用 AWS Kinesis 与 Lambda 处理流数据

Maksim Pecherskiy

Data Engineer

上一课

StkGt.png

使用 AWS Kinesis 与 Lambda 处理流数据

本课将学到

Zb4Va.png

使用 AWS Kinesis 与 Lambda 处理流数据

通用日志格式

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
使用 AWS Kinesis 与 Lambda 处理流数据

读取 S3 文件

import pandas as pd

结果:

Response:
{
  "errorMessage": "Unable to import module 'lambda_function': No module named 'pandas'",
  "errorType": "Runtime.ImportModuleError"
}
使用 AWS Kinesis 与 Lambda 处理流数据

更新处理器

#lambda_function.py
import json, boto3, pandas as pd
# 初始化客户端
...

SPEED_ALERT_THRESHOLD = 45 ALERT_PHONE_NUMBER = "+1234567890"
# 从写入记录获取数据框的辅助函数 def get_new_data(event): pd.read_csv()... ... return data
# Lambda 函数处理器 def record_created_handler(event, context): data = get_new_data(event) ... sns.publish() ...
使用 AWS Kinesis 与 Lambda 处理流数据

get_new_data()

def get_new_data(event):
    # 创建用于存储新对象键的列表。
    written_objects = []

    # 遍历每个 S3 事件记录。
    for record in event['Records']:

        # 获取要检查的变量
        event_name = record['eventName']
        bucket_name = record['s3']['bucket']['name']
        obj_key = record['s3']['object']['key']
使用 AWS Kinesis 与 Lambda 处理流数据

get_new_data()

def get_new_data(event):
        ...
        # 验证事件来自 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)
# 将新记录合并为单个数据框。 return pd.concat(written_objects)
使用 AWS Kinesis 与 Lambda 处理流数据

record_created_handler()

SPEED_ALERT_THRESHOLD = 45
ALERT_PHONE_NUMBER = "+1234567890"
...
def record_created_handler(event, context):
    # 调用辅助方法
    data = get_new_data(event)

## 获取最高速度 top_speeds = data.groupby(['vin'])['speed'].max().reset_index()
## 获取超过 45 的速度 too_fast = top_speeds.loc[top_speeds.speed > SPEED_ALERT_THRESHOLD, :]
使用 AWS Kinesis 与 Lambda 处理流数据

record_created_handler()

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

def record_created_handler(event, context):
    ...
    ## 发送短信
    sns.publish(PhoneNumber=ALERT_PHONE_NUMBER,
        Message="Speeding Alert \n" + too_fast.to_string())

## 目前不输出到任何地方,但需要返回点东西。 totals = data.groupby(['vin'])['speed'].max().reset_index() return totals.to_csv(sep=" ", index=False)
使用 AWS Kinesis 与 Lambda 处理流数据

测试 Lambda 函数

jBpMo.png

KodWB.png

使用 AWS Kinesis 与 Lambda 处理流数据

添加环境变量

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):
    ...
使用 AWS Kinesis 与 Lambda 处理流数据

添加环境变量

7jD6x.png

使用 AWS Kinesis 与 Lambda 处理流数据

添加环境变量

eDRTX.png

使用 AWS Kinesis 与 Lambda 处理流数据

yKjWC.png

使用 AWS Kinesis 与 Lambda 处理流数据

添加触发器

gLmbZD.png

使用 AWS Kinesis 与 Lambda 处理流数据

添加触发器

7F08p.png

使用 AWS Kinesis 与 Lambda 处理流数据

添加触发器

bPUtX.png

使用 AWS Kinesis 与 Lambda 处理流数据

复习

ZAbfT.png

使用 AWS Kinesis 与 Lambda 处理流数据

复习

m9jlv.png

使用 AWS Kinesis 与 Lambda 处理流数据

复习

GTlyD.png

使用 AWS Kinesis 与 Lambda 处理流数据

复习

XddOW.png

使用 AWS Kinesis 与 Lambda 处理流数据

Passons à la pratique !

使用 AWS Kinesis 与 Lambda 处理流数据

Preparing Video For Download...