轉換型 Lambda

使用 AWS Kinesis 與 Lambda 的串流資料

Maksim Pecherskiy

Data Engineer

第 1 章中…

2020-07-15_06-18.png

使用 AWS Kinesis 與 Lambda 的串流資料

第 2 章中…

2020-07-15_06-20.png

使用 AWS Kinesis 與 Lambda 的串流資料

本章將…

2020-07-15_06-27.png

使用 AWS Kinesis 與 Lambda 的串流資料

先前作法

2020-07-14_07-42.png

使用 AWS Kinesis 與 Lambda 的串流資料

在 S3 一次處理 vs Lambda 轉換

在 S3 中一次處理

  • 使用在 S3 物件寫入時觸發的 lambda 函式
  • 資料轉換延遲較長
  • 清理前先儲存原始資料
  • 需將 Firehose 目的地設為 S3

透過 Lambda 轉換處理

  • 使用在 Firehose 串流中途觸發的 lambda 函式
  • 立即完成轉換
  • 僅儲存清理後資料
  • Firehose 可設定為其他目的地
使用 AWS Kinesis 與 Lambda 的串流資料

傳入資料

record_id timestamp vin lon lat speed
939ed1d1-1740-420c-8906-445278573c7f 4:25:06.000 4FTEX4944AK844294 106.9447146 -6.3385652 25
f29a5b3d-d0fa-43c0-9e1a-e2a5cdb8be7a 8:10:47.000 3FTEX1G5XAK844393 108.580681 34.79925 37
ff8e7131-408d-463b-8d07-d016419b0656 20:26:44.000 2LAXX1C8XAK844292 114.392392 36.097577 90
使用 AWS Kinesis 與 Lambda 的串流資料

轉換型 Lambda

2020-07-15_06-38.png

使用 AWS Kinesis 與 Lambda 的串流資料

範例事件

{
  "invocationId": "invocationIdExample",
  "deliveryStreamArn": "arn:aws:firehose:us-east-1:458912630:deliverystream/gps-delivery-stream",
  "region": "us-east-1",
  "records": [
    {

"recordId": "49546986683135544286507457936321625675700192471156785154",
"approximateArrivalTimestamp": 1495072949453,
"data": "NjQuMjQyLjg4LjEwIC0gLSBbMDcvTWFyLzIwMDQ6MTY6MTA6MDIgLTA4MDBdICJHRVQgL21haWxtYW4vbGlzdGluZm8vaHNkaXZpc2lvbiBIVFRQLzEuMSIgMjAwIDYyOTE=="
} ] }
使用 AWS Kinesis 與 Lambda 的串流資料

Base64

2020-07-19_10-54.png

使用 AWS Kinesis 與 Lambda 的串流資料

處理資料

import base64
from datetime import datetime as dt

def convert_timestamp(record_time_val):
# 取得今日日期字串 today = dt.today().strftime("%Y-%m-%d")
# 用今日日期結合記錄時間,建立 datetime 物件 new_ts = dt.strptime(f"{today} {record_time_val}", "%Y-%m-%d %H:%M:%S.%f")
# 將 datetime 物件轉為易讀字串 return new_ts.strftime("%Y-%m-%dT%H:%M:%S")
使用 AWS Kinesis 與 Lambda 的串流資料

處理資料

def lambda_handler(event, context):
    output = []

for record in event['records']: payload = base64.b64decode(record['data'])
payload = payload.decode()
payload = payload.split(" ")
payload[1] = convert_timestamp(payload[1])
使用 AWS Kinesis 與 Lambda 的串流資料

處理資料

    for record in event['records']:
        ...
        payload = " ".join(payload)
        payload_enc = base64.b64encode(payload.encode())

output.append({
'recordId': record['recordId'],
'result': 'Ok',
'data': payload_enc
})
return {'records': output}
使用 AWS Kinesis 與 Lambda 的串流資料

重點回顧

import base64
def lambda_handler(event, context):

output = [] for record in event['records']: # 逐筆處理記錄
payload = base64.b64decode(record['data']).decode() # 解碼載荷
# 進行修改
payload_enc = base64.b64encode(payload.encode()) # 重新編碼
output.append({ # 放入字典 'recordId': record['recordId'], 'result': 'Ok', 'data': payload_enc, })
return {'records': output}
使用 AWS Kinesis 與 Lambda 的串流資料

在 AWS 建立 Lambda

使用 AWS Kinesis 與 Lambda 的串流資料

建立 lambda 的指令稿

使用 AWS Kinesis 與 Lambda 的串流資料

一起來練習吧!

使用 AWS Kinesis 與 Lambda 的串流資料

Preparing Video For Download...