高级阅读约 3 分钟

DynamoDB Streams 完整指南(含示例)

DynamoDB Streams 是一个变更数据捕获日志:一张表上的每一次插入、更新和 删除都被有序地捕获,形成一个你可以对其做出反应的记录流。 它是你在不轮询一张表的情况下把它变成事件源的方式。

在审计日志的场景里,你想在一个敏感事件落地的瞬间做出反应 —— 在有人导出一张发票或授予一个管理员角色时触发一次告警 —— 而 无需定时扫描那张表。Streams 就是那件事的推送侧。

DynamoDB Streams 是如何工作的?

DynamoDB Streams 把一张表上的每一次插入、更新和删除都捕获为一个按时间排序、经过去重的记录日志,保留最长 24 小时。你用 StreamViewType(键、新映像、旧映像,或两者)选择每条记录携带什么,然后用一个 Lambda 触发器消费这个流,从而在不轮询的情况下对项变更做出反应。

  • Streams 把项级变更捕获为一个按时间排序、经过去重的日志, 保留最长 24 小时
  • 你选择每条记录携带什么,通过 StreamViewType:仅键、新 映像、旧映像,或旧与新两者。
  • 记录按项有序 —— 对一个项的多次变更按它们被 写入的顺序到达 —— 而一个流像一样被分片。
  • 原生消费者是 Lambda —— 一个针对每批新记录运行的触发器, 而 Kinesis Data Streams 是用于更丰富扇出的替代方案。

问题:在不轮询的情况下做出反应

你需要"在一个 role.granted 事件被写入时告警我"。天真的做法是 一个每分钟扫描新事件的定时作业 —— 它每次都读取整个 近期分区,耗费容量,而且总是至少晚一分钟。

你真正想要的是一次推送:DynamoDB 在一个项变更的那一刻 告诉你。那正是 Streams 所提供的,变更记录会被送到 你的代码,而不是你去追猎它。

Streams 是如何运作的

按 AWS 文档所述,DynamoDB Streams 把变更保留为一个经过去重、 按时间排序的日志,最长 24 小时,并原生集成 Lambda (DynamoDB 的变更数据捕获)。 每条记录描述一次项级的修改。

当你启用一个流时,你挑选一个 StreamViewType,它控制每条记录 携带被变更项的多少内容:

StreamViewTypeeach record contains
KEYS_ONLYonly the key attributes of the changed item
NEW_IMAGEthe entire item as it looks after the change
OLD_IMAGEthe entire item as it looked before the change
NEW_AND_OLD_IMAGESboth the before and after images

记录按项有序 —— 对单个项的多次变更按它们被 写入的顺序出现 —— 而流沿着与表相同的分区结构被分片。 保留期是 24 小时 —— Streams 是一个 反应缓冲区,不是一份永久历史。要保留可持久的历史,你就存储事件 本身(我们的审计日志表恰恰就是这个)。

原生消费者是一个 Lambda 触发器:DynamoDB 在一批新流记录到达时, 用它们调用你的函数。

LambdaStream"DynamoDB"AppLambdaStream"DynamoDB"App"Put EVENT role.granted""变更记录(NEW_IMAGE)""一批记录""若 action 敏感 → 告警"

一个完整示例:对敏感审计事件告警

审计日志表配上一个 NEW_IMAGE 的流,这样每条记录都携带 完整的新事件。一个 Lambda 消费这一批,只转发那些 要紧的记录:

stream record (NEW_IMAGE)consumer action
TENANT#acmeEVENT#…#a2action=invoice.exportsend to SIEM
TENANT#globex EVENT#…#b9 action=role.grantedpage on-call
TENANT#acmeEVENT#…#a1action=login.successignore

这个函数从不触及那张表 —— 它纯粹对流交给它的东西 做出反应。无轮询、无扫描,而且告警在写入后几秒内 触发。流记录按项有序,所以对同一个事件项的连续变更 按它们被写入的顺序到达。

这也是维护一份下游副本的标准方式:一个流 消费者可以把每个事件投影进 OpenSearch 以做全文审计搜索,或者 聚合计数 —— 全都从同一个变更日志派生。

在 DynoTable 中操作

在你接线一个流消费者之前,你需要知道你的 Lambda 将要接收的 那个项的确切形态 —— 有哪些属性、嵌套的映射和列表长什么 样、一条 NEW_IMAGE 记录实际会包含什么。

要在纯 JSON 与一条流记录所用的属性值形态之间转换一个样本 项,DynamoDB JSON 转换器会在 你的浏览器里完成它。而在 DynoTable 里,你可以检视完整的项 —— 包括它的 DynamoDB-JSON 形态 —— 从而针对真实数据建模那条 NEW_IMAGE 记录,而不是 猜测字段的形态。

在 DynoTable 中检视一个审计事件项,以建模其 Lambda 消费者将接收的 NEW_IMAGE 流记录。
在 DynoTable 中检视一个审计事件项,以建模其 Lambda 消费者将接收的 NEW_IMAGE 流记录。

如果你在本地测试一个消费者,就针对 DynamoDB Local 运行那张表,并 以同样的方式检视它 —— 参见 连接到 DynamoDB Local

陷阱与后续步骤

  • 24 小时不是一个积压队列。 如果你的消费者宕机一天,记录就会 过期消失。Streams 用于近实时反应,而不是可持久的重放 —— 为了历史,请保留事件本身。
  • 挑选你需要的最小 StreamViewType NEW_AND_OLD_IMAGES 会让 载荷翻倍;如果你只需要键去重新读取那个项,KEYS_ONLY 更廉价。
  • 有序性是按项的,不是按分区键或全局的。 DynamoDB 只对同一个项的 连续变更保证顺序;跨不同项之间没有顺序保证, 即便在同一个分区键内也如此。
  • TTL 删除会以流记录的形式出现,带有系统属性标记,这 正是你归档过期项的方式 —— 参见 DynamoDB TTL

Streams 把审计日志变成一个事件源。下一个运维关注点是 一个项生命的另一端 —— 用 DynamoDB TTL 自动过期旧事件。

下载 DynoTable,在你写一行 Lambda 代码之前,检视你的流 消费者将要接收的那个确切的项形态。

更新于