Lambda递归Promise查询DynamoDB过慢问题及优化咨询
优化DynamoDB分片GSI查询性能的方案
问题背景
我采用Alex Debrie在其DynamoDB书籍中提到的分片模式,从表中获取最新创建的条目。写入新条目时,会将数据同步到GSI,结构如下:
- GSI1PK:截断后的时间戳#[0-9]的随机分片号(例如
2023-02-09T00:00:00.000Z#8) - GSI1SK:条目的唯一ID
查询时通过三个参数获取最新条目:
- Date:当前日期
- Limit:需要获取的条目总数
- Days:回溯查询的天数
按照书籍建议,我用带Promise的递归函数实现查询,但Lambda函数运行极慢:当近期条目数量少且分散在各分片时,比如要拉取过去7天的100条但实际不足100条,函数需要执行70次Query(7天×10个分片),耗时约10秒;而近期条目多的时候仅需1秒左右。
环境信息
- 单条条目大小约400字节
- DynamoDB表使用按需容量模式
- Lambda配置1536MB内存
- Node.js 16.x
当前实现代码
const getQueryParams = (createdAt, shard, limit) => { const params = { TableName : "table", IndexName: 'GSI1', KeyConditionExpression: "#gsi1pk = :gsi1pk", ExpressionAttributeNames: { "#gsi1pk": 'GSI1PK' }, ExpressionAttributeValues: { ":gsi1pk": `${truncateTimestamp(timestamp).toISOString()}#${shard}` //e.g 2023-02-09T00:00:00.000Z#8 }, ScanIndexForward: false, Limit: limit }; return params; } const getItems = async => { const items = [] const number_of_days = 3; const getLatestItems = async ({ createdAt = new Date(), limit = 100, days = 0, shard = 0 }) => { const query_params = getQueryParams(createdAt, shard, limit); let max_items_to_fetch = limit; return dynamoDb.query(query_params).then( (data) => { // process data. if (data.Items) { data.Items.forEach((item) => { if (items.length < limit) { items.push(item); } }) max_items_to_fetch = limit - data.Items.length; } if (items.length >= limit) { return items; } if (shard < 9) { let params = { createdAt: new Date(createdAt.setDate(createdAt.getDate())), limit: max_items_to_fetch, days: days, shard: shard + 1, } return getLatestItems(params); } else if (days < number_of_days) { let params = { createdAt: new Date(createdAt.setDate(createdAt.getDate() - 1)), limit: max_items_to_fetch, days: days + 1, shard: 0, } return getLatestItems(params); } return items; }, (error) => { throw new Error('Error getting all recent itmems') } ); } return getLatestItems({}); }; export const main = async (event) => { const start = Date.now(); const itemPromises = getItems(); const res = await Promise.all([itemPromises]); const end = Date.now(); console.log(`Execution time: ${end - start} ms`); };
优化方案
1. 并行查询分片,消除串行等待
当前递归逻辑是逐个串行遍历分片和日期,这是性能瓶颈的核心。改为对同一天的所有分片发起并行查询,待当天所有分片结果返回后,再处理前一天的分片,能大幅减少等待时间。
优化后的核心代码:
const getItems = async (targetLimit = 100, daysToBacktrack = 7) => { const items = []; let currentDate = new Date(); currentDate.setHours(0, 0, 0, 0); // 与GSI1PK的时间截断逻辑保持一致 // 处理单天的所有分片查询 const processSingleDay = async (date) => { if (items.length >= targetLimit) return; // 并行发起当天10个分片的查询 const shardQueryPromises = Array.from({ length: 10 }, (_, shard) => { const queryParams = getQueryParams(date, shard, targetLimit - items.length); return dynamoDb.query(queryParams).promise(); }); // 等待所有分片查询完成,合并结果 const queryResults = await Promise.all(shardQueryPromises); queryResults.forEach(result => { if (result.Items) { result.Items.forEach(item => { if (items.length < targetLimit) { items.push(item); } }); } }); }; // 按日期倒序遍历,处理完当天再切换到前一天 for (let i = 0; i < daysToBacktrack; i++) { await processSingleDay(currentDate); if (items.length >= targetLimit) break; // 切换到前一天 currentDate = new Date(currentDate.setDate(currentDate.getDate() - 1)); } // 对合并后的结果按创建时间重新排序(各分片结果是局部倒序,整体需统一排序) items.sort((a, b) => new Date(b.createdAt) - new Date(a.createdAt)); // 确保返回数量不超过目标值 return items.slice(0, targetLimit); };
2. 修复原代码中的逻辑错误
getItems函数定义语法错误:async =>需改为async () =>- 剩余需获取条目数计算错误:
max_items_to_fetch = limit - data.Items.length应改为limit - items.length,避免重复计数 - Date对象副作用问题:直接调用
createdAt.setDate()会修改原对象,需创建新的Date实例避免逻辑混乱
3. 按需调整分片粒度(可选)
如果写入量允许,可将分片数量从10个减少到5个,减少空分片的无效查询次数。但需注意分片数量过少可能引发写入热点,需根据业务写入量权衡。
4. 减少数据传输量
如果只需要条目的部分字段,在getQueryParams中添加ProjectionExpression,仅返回所需字段,降低网络传输开销:
const getQueryParams = (createdAt, shard, limit) => { const params = { TableName : "table", IndexName: 'GSI1', KeyConditionExpression: "#gsi1pk = :gsi1pk", ExpressionAttributeNames: { "#gsi1pk": 'GSI1PK' }, ExpressionAttributeValues: { ":gsi1pk": `${truncateTimestamp(createdAt).toISOString()}#${shard}` }, ScanIndexForward: false, Limit: limit, ProjectionExpression: "id, createdAt, content" // 替换为实际需要的字段 }; return params; }
5. 添加查询缓存(可选)
如果查询时效性要求不高,可在Lambda中引入缓存(如ElastiCache或直接用DynamoDB存储近期查询结果),避免重复查询相同日期范围的数据。
内容的提问来源于stack exchange,提问作者mvp
相关产品推荐
相关产品推荐

