You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 11:25:30