AWS Kinesis と Lambda で学ぶストリーミングデータ
Maksim Pecherskiy
Data Engineer


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
import pandas as pd
results in:
Response:
{
"errorMessage": "Unable to import module 'lambda_function': No module named 'pandas'",
"errorType": "Runtime.ImportModuleError"
}
#lambda_function.py import json, boto3, pandas as pd # Initialize clients ...SPEED_ALERT_THRESHOLD = 45 ALERT_PHONE_NUMBER = "+1234567890"# 新規レコードから DataFrame を取得するヘルパー def get_new_data(event): pd.read_csv()... ... return data# Lambda ハンドラー def record_created_handler(event, context): data = get_new_data(event) ... sns.publish() ...
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']
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)
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, :]
SPEED_ALERT_THRESHOLD = 45 ALERT_PHONE_NUMBER = "+1234567890" ... def record_created_handler(event, context): ... ## SMS 送信 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)


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 で学ぶストリーミングデータ