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

AWS Lambda中Promise.all本地正常运行 部署后callback返回null问题排查

AWS Lambda 并发请求提前终止异常排查

问题描述

该Lambda函数的核心逻辑是通过批量API调用获取多组数据,再通过FiveTran同步到Snowflake。当前异常表现为:代码使用Promise.all()处理并发请求,本地环境运行完全符合预期,但部署到AWS后callback返回null。通过Cloudwatch日志排查发现,代码执行到lib.allAPICalls函数内的Promise.all()逻辑后就直接停止退出,函数超时配置为60秒,但实际部署后约6秒就返回null,无法定位问题根因。

问题代码

'use strict'

import { 
    Callback,
    Context,
    Handler,
} from "aws-lambda"

import { ServiceException } from '@destpet/dp-validation'
import { Logger } from '@destpet/dp-logger'
import fetch from 'isomorphic-fetch'
const trace = Logger('trace')
const debug = Logger('debug')
const warn = Logger('warn')
const error = Logger('error')
import data from "./data"
import { createBuilderStatusReporter, forEachTrailingCommentRange, getDefaultLibFilePath } from "typescript"

const ENV:string = process.env['NODE_ENV']
const lambda:string = 'dp_PredictHQ_etl_v1' 
const version:string = '1.0.0'

const CORS_ALL_ORIGINS = {'Access-Control-Allow-Origin': '*'}

export interface FiveTranRequest {
    agent: string;
    state: {transactionsCursor: string};
    secrets: {apiToken: string};
}

export interface FiveTranResponse {
    state: {transactionsCursor: string};
    insert: {transactions:Array<object>};
    delete?: {transactions:Array<object>};
    schema?: {transactions: {primary_key: Array<string>}};
    hasMore: boolean;
    test: any;
}

export interface FiveTranError {
    errorMessage: string;
    errorType: string;
    stackTrace?: Array<Array<string>>;
}

export interface LocationInfo {
    Location_ID: string;
    Latitude: string;
    Longitude: string;
}

export interface AllAPICallsReturn {
    responsePromiseArray: any,
    locations: Array<number>
}

const lib = {

    /**
     * Process service request
     */
     processRequest: async (
            event: FiveTranRequest,
            resp: FiveTranResponse, 
            callback:Callback
        ): Promise<void> => {

        try {
            trace('just inside of processrequest',`${lambda}: processRequest`, null, resp)
            const endDate = new Date
            endDate.setDate(endDate.getUTCDate() + 90)
            const dateArray = lib.createDateRange(endDate)
            
            const locationArray = data;
            
            let transactionsArray = []
            
            let locationArrayPointer = 0
            let pointerIncrement = 10
            
            while (locationArrayPointer <= locationArray.length) {
                let templocationArray = locationArray.slice(locationArrayPointer,locationArrayPointer + pointerIncrement)
                let apiDataFetch = await lib.allAPICalls(dateArray, templocationArray)
                let apiDataParse = lib.parseAllData(apiDataFetch.fetchPromiseArray, apiDataFetch.locations)
                transactionsArray = transactionsArray.concat(apiDataParse)
                locationArrayPointer += pointerIncrement;
            }
            
            resp['state']['transactionsCursor'] = '2021-10-04'
            resp['schema']['transactions'] = {primary_key: ['id']}
            resp['hasMore'] = false
            resp['insert']['transactions'] = transactionsArray
            resp['test'] = transactionsArray
            
            
            // trace('right before callback',`${lambda}: processRequest`, null, resp)
            callback(null, resp)

        } catch (e:any) {
            error('Unable to process request', `${lambda}:processRequest`, null, e)
            console.log(e)

            let stackTrace = (e instanceof Error && e.stack) ? [e.stack.split('\n')] : [['No stack trace']]

            const errResponse: FiveTranError = {
                errorMessage: (e instanceof ServiceException) ? e.message : e.toString(),
                errorType: 'DP_FiveTran_API_Lambda_Error',
                stackTrace
            }

            const errorCode = (e instanceof ServiceException) ? e.errorCode : 500

            callback(null, resp)
        }
    },

    createDateRange : ( endDate : Date) :Array<string> => {
       let dateArray = [];

       let datePointer = new Date()
       datePointer.setUTCFullYear(2017)
       datePointer.setUTCMonth(7)
       datePointer.setUTCDate(1)

       while(endDate > datePointer) {
           let internalArray = []
           let startInterval = new Date(datePointer.getUTCFullYear(), datePointer.getUTCMonth(), 1)
           let endInterval = new Date(datePointer.getUTCFullYear(), datePointer.getUTCMonth() + 2, 0)
           internalArray.push(startInterval.toISOString().split('T')[0])
           internalArray.push(endInterval.toISOString().split('T')[0])
           dateArray.push(internalArray)
           datePointer = new Date(datePointer.getUTCFullYear(), datePointer.getUTCMonth() + 2, 1)
       }

       dateArray[dateArray.length - 1][1] = endDate.toISOString().split('T')[0]

       return dateArray
    },

    APICall : async (date: string, location: LocationInfo) => {
        
        return new Promise((resolve, reject) => {
            const dateS = date[0]
            const dateE = date[1]
            const data = {
                "active": {
                    "gte": dateS,
                    "lte": dateE
                },
                'phq_rank_school_holidays': true,
                'phq_rank_public_holidays': true,
                'phq_rank_observances': true,
                'phq_rank_academic_holiday': true,
                'phq_attendance_academic_social': true,
    
                'location': {
                    'geo': {
                        'lat': Number(location.Latitude),
                        'lon': Number(location.Longitude),
                        'radius': '12km'
                    }        
                },
                "phq_attendance_sports": {
                    "stats": [
                        "sum"
                    ]
                }
            }
    
            const requestBody = {
                headers: {
                "Authorization": "redacted",
                "Accept": "application/json"
                },
                method: 'POST',
                body: JSON.stringify(data),
                'Content-type': 'application/json' 
            }

            fetch('https://api.predicthq.com/v1/features', requestBody)
                .then(r => {
                    return r.json();
                })
                .then(response => {
                    resolve(response);
                })
                .catch(err => {
                    reject(new ServiceException('APICall promise rejected', 500))
                    error('in APICall','',null,error)
                })
        })
    },
    
    allAPICalls : async (dateArray : Array<string>, locationsArray : Array<LocationInfo>): Promise<any> => {
        const promiseArray: Array<Promise<any>> = [];
        const locations: Array<number> = [];
        
        locationsArray.forEach( location => {
            dateArray.forEach( date => {
                promiseArray.push(lib.APICall(date, location))
                locations.push(Number(location.Location_ID))
            })
        })
        
        let fetchPromiseArray, responsePromiseArray
        try {
            fetchPromiseArray = await Promise.all(promiseArray).catch(e => {error('','',null,e)})
        } catch (err) {
            error('error in all api calls',`${lambda}: processRequest`, null, err)
            console.log(err)
        }
        return {
            fetchPromiseArray,
            locations
        }
    },

    dataParse : (dataObject) => {
        const parsedData = {};
        const recursiveParse = ( recursiveObject = dataObject, fieldTitle = '' ) => {
           for (const property in recursiveObject) {
                if (property !== 'phq_attendance_academic_social') {
                    let dataFieldTitle
                    if (fieldTitle === '') {
                        dataFieldTitle = property
                    } else {
                        dataFieldTitle = fieldTitle + '_' + property
                    }
                    if (typeof recursiveObject[property] !== 'object') {
                        parsedData[dataFieldTitle] = recursiveObject[property]
                    } else {
                        recursiveParse(recursiveObject[property], dataFieldTitle)
                    }
               }
           } 
        }
        recursiveParse();
        return parsedData
    },

    parseAllData : ( dataArray, locationsArray) => {
        const transactionsArray = [];
        let id = 0;
        dataArray.forEach( (resultArray, index) => {
            resultArray.results.forEach( individualResult => {
                let transaction = lib.dataParse(individualResult)
                transaction['Location'] = locationsArray[index]
                transaction['id'] = id;
                transactionsArray.push(transaction)
                id += 1;
            })
        })
        return transactionsArray;
    }
}

/**
 * Lambda service handler
 * 
 * @param event 
 * @param context 
 * @param callback 
 */
 export const handler: Handler<any, void> = async (
        event: FiveTranRequest,
        context: Context, 
        callback: Callback
    ): Promise<void> => {

    let resp: FiveTranResponse = {
        state: {transactionsCursor: null},
        schema: {transactions: null},
        insert: {transactions: null},
        hasMore: false,
        test: ''
    }

    if (event['detail-type'] === 'warm') {
        trace('...Warm...')
        callback(null, null)
    } else {
        debug(`Executing service handler`, `${lambda}:handler`)
        lib.processRequest(event, resp, callback).catch(err => {console.log(`Error: ${err}`)})
    }
}

根因与解决方案

根因说明

AWS Lambda的异步handler函数在执行完成后会立即冻结执行上下文,不会等待未被await的异步任务执行完成。本地开发环境不会主动回收执行上下文,因此异步逻辑可以正常执行完成,而部署到AWS后,lib.processRequest是异步函数,未加await的情况下handler会直接执行完成,导致异步逻辑被中断,提前返回null。

修复方案

在handler函数中调用lib.processRequest的位置添加await关键字即可,修改后的handler代码如下:

export const handler: Handler<any, void> = async (
        event: FiveTranRequest,
        context: Context, 
        callback: Callback
    ): Promise<void> => {

    let resp: FiveTranResponse = {
        state: {transactionsCursor: null},
        schema: {transactions: null},
        insert: {transactions: null},
        hasMore: false,
        test: ''
    }

    if (event['detail-type'] === 'warm') {
        trace('...Warm...')
        callback(null, null)
    } else {
        debug(`Executing service handler`, `${lambda}:handler`)
        // 新增await关键字等待异步任务执行完成
        await lib.processRequest(event, resp, callback).catch(err => {console.log(`Error: ${err}`)})
    }
}

内容的提问来源于stack exchange,提问作者Nicholas Brush

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 22:36:05