DynamoDB Streams是DynamoDB提供的变更数据捕获机制,每当表中的数据发生新增、修改或删除时,相关变更记录会被写入流中,供下游服务消费。在Nodejs生态里,消费这些流数据主要有两条路:一是配置Lambda触发器让AWS托管整个消费过程,二是用AWS SDK自己写一个Consumer长轮询拉取数据。两条路线各有优劣,选错了场景往往会导致数据丢失、重复处理或者成本失控。这篇文章把两条路线都拆开讲清楚,并附上完整的代码示例。

一、DynamoDB Streams的工作原理与前置配置
先理解流的底层模型。DynamoDB Streams与Kinesis在概念上高度相似,每张表开启流之后,数据按分区(Partition)组织成多个分片(Shard),每个分片内部维护一条有序的记录序列。消费者通过分片迭代器(ShardIterator)标记读取位置,逐批拉取记录。理解这一点非常重要,因为无论用Lambda还是自建Consumer,底层都是同一套机制。
开启流需要在表级别配置StreamSpecification,可以指定四种视图类型:NEW_IMAGE(新值)、OLD_IMAGE(旧值)、NEW_AND_OLD_IMAGES(新旧值都包含)、KEYS_ONLY(仅主键)。视图类型的选择直接影响存储成本和下游处理逻辑,比如你要做字段级别的diff对比,就必须选NEW_AND_OLD_IMAGES,否则拿不到变更前的数据。用Nodejs开启流的代码如下:
const { DynamoDBClient, UpdateTableCommand } = require('@aws-sdk/client-dynamodb');
const client = new DynamoDBClient({ region: 'ap-northeast-1' });
await client.send(new UpdateTableCommand({
TableName: 'Orders',
StreamSpecification: {
StreamEnabled: true,
StreamViewType: 'NEW_AND_OLD_IMAGES'
}
}));还有一个容易忽略的细节:流记录的保留期固定为24小时,不可调整。如果消费者宕机超过24小时没有拉取数据,未处理的记录就会被丢弃,这一点在规划重试策略和容灾方案时必须考虑进去。
二、Lambda触发器:托管消费的首选方案
绝大多数场景下,Lambda是消费DynamoDB Streams最省心的方式。你不需要管理迭代器、不需要处理分片的分裂与合并,AWS会自动把流记录打包成批量事件调用你的函数。在Nodejs里写一个基础的Lambda消费者非常简单:
exports.handler = async (event) => {
for (const record of event.Records) {
// 判断变更类型:INSERT / MODIFY / REMOVE
console.log('eventName:', record.eventName);
// 新值在 dynamodb.NewImage 中,旧值在 dynamodb.OldImage 中
const newImage = record.dynamodb.NewImage;
const oldImage = record.dynamodb.OldImage;
// 注意:字段值是 DynamoDB 的属性格式,需要手动解包
if (newImage && newImage.status) {
console.log('status:', newImage.status.S);
}
}
return { batchItemFailures: [] };
};这里有两个关键点。第一,流记录中的字段值是DynamoDB原生属性格式,字符串对应S、数字对应N,直接当普通对象访问会拿到undefined,建议封装一个解包工具函数或者使用unmarshall工具转换。第二,返回值中的batchItemFailures配合函数配置里的FunctionResponseTypes: ['ReportBatchItemFailures']使用,可以做到记录级别的失败重试,而不是整批重试,这对降低重复处理量非常关键。
const { unmarshall } = require('@aws-sdk/util-dynamodb');
exports.handler = async (event) => {
const failedIds = [];
for (const record of event.Records) {
try {
const item = record.dynamodb.NewImage
? unmarshall(record.dynamodb.NewImage)
: null;
await processRecord(item, record.eventName);
} catch (err) {
// 记录失败项的序列号,Lambda只重试失败的记录
failedIds.push({ itemIdentifier: record.dynamodb.SequenceNumber });
}
}
return { batchItemFailures: failedIds };
};还需要注意并发与乱序问题。Lambda对同一个分片保证串行处理,但不同分片之间是并行的,且并行度默认与表的分区数挂钩。写入量大的热点表分区多,Lambda并发也会相应升高,如果你的下游处理逻辑对并发敏感(比如写数据库有连接数限制),要提前评估。另外,分片内的修改记录是有序的,但同一主键的多次变更可能落在同一批次内,业务逻辑要能处理中间态。
三、自建Consumer:用AWS SDK手动拉取流数据
当消费逻辑需要运行在自有服务器、容器集群,或者需要对拉取节奏做精细控制时,就要自己写Consumer了。核心流程分四步:列出分片、获取迭代器、循环调用GetRecords、处理分片状态。下面是一个可运行的骨架代码:
const {
DynamoDBStreamsClient,
ListStreamsCommand,
DescribeStreamCommand,
GetShardIteratorCommand,
GetRecordsCommand
} = require('@aws-sdk/client-dynamodb-streams');
const streams = new DynamoDBStreamsClient({ region: 'ap-northeast-1' });
async function consumeShard(streamArn, shardId) {
// 获取迭代器,TRIM_HORIZON表示从最旧的记录开始读
const { ShardIterator } = await streams.send(new GetShardIteratorCommand({
StreamArn: streamArn,
ShardId: shardId,
ShardIteratorType: 'TRIM_HORIZON'
}));
let iterator = ShardIterator;
while (true) {
const res = await streams.send(new GetRecordsCommand({
ShardIterator: iterator,
Limit: 1000
}));
if (res.Records && res.Records.length > 0) {
for (const r of res.Records) {
console.log(r.eventName, r.dynamodb.Keys);
}
}
if (res.NextShardIterator) {
iterator = res.NextShardIterator;
} else {
// 分片已关闭,需要处理子分片
break;
}
// 空轮询时适当休眠,避免打爆API配额
if (!res.Records || res.Records.length === 0) {
await new Promise(r => setTimeout(r, 1000));
}
}
}自建方案有几个必须处理的坑。首先是迭代器有效期只有15分钟,但每次GetRecords会返回新的迭代器,只要持续轮询就没问题,程序重启后则需要用AT_SEQUENCE_NUMBER或AFTER_SEQUENCE_NUMBER结合自己持久化的位点来恢复。其次,GetRecords单次最多返回1000条且数据量不超过1MB,被截断时要继续用返回的迭代器拉取,不能丢弃。最后,分片会随分区变化而分裂,父分片关闭后要通过DescribeStream发现子分片并继续消费,否则会静默漏数据。
如果没有特殊部署需求,其实更推荐的做法是把自建Consumer逻辑放到ECS或K8s里跑,配合checkpoint存储(比如把序列号写到Redis或DynamoDB自身),实现至少一次投递语义。切记流消费是至少一次而非恰好一次,下游逻辑必须做幂等设计,最简单的方式是用序列号或主键加版本号做去重。
四、生产环境的调优与监控要点
性能层面,Lambda方案主要调三个参数:BatchSize(单批记录数,最大1000)、MaximumBatchingWindow(攒批窗口,最大300秒)和并发保留配置。数据量小但要求低延迟,就缩小窗口;数据量大且下游吞吐有限,就加大BatchSize并开启失败记录上报。自建方案则要控制好轮询节奏,GetRecords每个分片每秒最多5次调用,空轮询务必加退避。
监控方面重点盯三个指标:一是IteratorAge(迭代器年龄),它反映消费滞后程度,持续增长说明消费者跟不上下游写入速度;二是函数错误率与重试次数,配合DLQ(死信队列)兜底处理失败批次;三是分片数量变化,热点表分区扩容会带来分片分裂,自建Consumer要能动态响应。对于Lambda方案,建议开启Destination配置,把彻底处理失败的批次路由到SQS或SNS,方便后续人工介入或补偿。
总结一下选型建议:常规事件驱动场景直接用Lambda触发器,省心且弹性好;有特殊运行环境要求、需要严格控速或复杂位点管理时再考虑自建Consumer。无论哪种方案,幂等处理和失败兜底都是不可省略的环节,这两点做扎实了,流式管道才能在生产环境长期稳定运行。
DynamoDB StreamsNodejsLambda触发器修改时间:2026-09-10 01:38:41