事件源映射与流处理

使用 AWS Lambda 构建无服务器应用

Claudio Canales

Senior DevOps Engineer

事件源映射的位置

  • 队列与流常以批量方式轮询。
  • 事件源映射管理轮询与批处理行为。

事件源映射位置

使用 AWS Lambda 构建无服务器应用

什么是事件源映射?

  • 连接到函数的资源。
  • Lambda 轮询源以获取记录。
  • Lambda 以批次调用处理函数。

事件源映射流程

使用 AWS Lambda 构建无服务器应用

推送 vs 轮询(为何重要)

推送模型

  • 推送源立即发送事件。
  • 例如:上传到 Amazon S3。

轮询模型

  • 轮询源由 Lambda 主动检查。
  • Lambda 组批并调用处理函数。
  • 这会影响延迟与重试行为。

推送 vs 轮询 模型

使用 AWS Lambda 构建无服务器应用

队列 vs 流

队列(Amazon SQS)

  • 托管消息队列,类似收件箱。

流(DynamoDB Streams)

  • DynamoDB 表的变更日志。
  • 都以 Records 到达,但含义不同。

队列 vs 流 对比

使用 AWS Lambda 构建无服务器应用

批事件结构(简化)

{
  "Records": [{
    "messageId": "abc-123",
    "body": "{\"order_id\": \"A-42\"}"
  }]
}
  • 大多数映射事件以 Records 开头。
  • 每条记录有 messageIdbody 字符串。
  • 解析 body 得到自定义负载。
使用 AWS Lambda 构建无服务器应用

演示:处理每条记录

def lambda_handler(event, context):
    records = event.get("Records", [])
    for record in records:
        body = record.get("body", "")
        print("BODY:", body)
    return {"statusCode": 200}
  • 读取 Records,提供默认空列表。
  • 遍历每条记录,安全读取 body
  • 记录必要日志后返回。
使用 AWS Lambda 构建无服务器应用

演示:解析 JSON body

import json

def lambda_handler(event, context):
    record = event.get("Records", [])[0]
    payload = json.loads(record.get("body", "{}"))
    order_id = payload.get("order_id")
    print("ORDER_ID:", order_id)
    return {"statusCode": 200}
  • 导入 json,并以安全默认值读取 body
  • json.loads 解析为字典。
  • 提取所需字段;先校验。
使用 AWS Lambda 构建无服务器应用

批处理:批大小与批窗口

  • 批大小:每次调用包含的记录数。
  • 批窗口:Lambda 等待凑满批次的时长。
  • 批越大吞吐越高,但单次处理开销增大。
  • 对于 SQS,关注 ApproximateAgeOfOldestMessage 以监控积压。

批大小权衡图

使用 AWS Lambda 构建无服务器应用

批大小的权衡

  • 大批量更高效,但可能增加延迟。
  • 小批量缩短首次处理时间,但调用次数更多。
  • 依据工作负载选择。

批大小 vs 延迟 图表

使用 AWS Lambda 构建无服务器应用

通过并发扩展

  • 多个批次可并行处理。
  • 更高并发提升吞吐。
  • 但也会增加下游负载。

并发批处理

使用 AWS Lambda 构建无服务器应用

保护下游系统

  • 过高并发会压垮数据库或 API。
  • 通过限流控制负载。
  • 倾向于小而安全的工作单元。

下游负载保护

使用 AWS Lambda 构建无服务器应用

部分批失败

  • 一个批中可能含好与坏记录。
  • 若不支持部分失败,一条坏记录会重试整个批。
  • 部分失败仅重试失败项。

部分批失败流程

使用 AWS Lambda 构建无服务器应用

演示:报告部分失败(SQS)

def lambda_handler(event, context):
    failures = []
    for record in event.get("Records", []):
        try:
            process(record)
        except Exception:
            failures.append({"itemIdentifier": record["messageId"]})
    return {"batchItemFailures": failures}
  • batchItemFailures 中收集失败记录的 messageId
  • Lambda 仅重试这些项。
  • 成功处理的记录不会重发。
使用 AWS Lambda 构建无服务器应用

关键要点

  • 事件源映射轮询队列与流。
  • 它们以批次将 Records 传给处理函数。
  • 批大小与窗口在延迟与吞吐间取舍。
  • 部分失败避免重处理已成功项。

映射要点总结图

使用 AWS Lambda 构建无服务器应用

Ayo berlatih!

使用 AWS Lambda 构建无服务器应用

Preparing Video For Download...