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
相关产品推荐
相关产品推荐

