Python Lambda导出CloudWatch日志到S3协程相关报错排查
问题背景
- 开发目标:编写Python类型的AWS Lambda函数,实现单次调用即可自动批量将多个CloudWatch日志组的日志导出到S3存储桶。
- 同步写法报错:未使用async/await语法时,调用
CreateExportTask接口返回如下错误:
An error occurred (LimitExceededException) when calling the CreateExportTask operation: Resource limit exceeded.
- 异步写法报错:改用async/await语法编写后,触发新的运行时序列化错误,需要排查修复。
原问题代码
import json import boto3 import datetime import os import time import asyncio async def lambda_handler(event, context): client = boto3.client('logs') days = 7 no_of_days = int(days) region_name = "ap-southeast-1" destination_bucket = "abc" response = client.describe_log_groups( limit=50 ) print(response) x = list() for i in response['logGroups']: x.append(i['logGroupName']) # while "nextToken" in response: # response = client.describe_log_groups( # nextToken=response["nextToken"] # ) # for i in response['logGroups']: # x.append(i['logGroupName']) # print(x) current_datetime = datetime.datetime.now() startTime = current_datetime - datetime.timedelta(days=no_of_days) endTime = current_datetime - datetime.timedelta(days=0) fromDate = int(startTime.timestamp() * 1000) toDate = int(endTime.timestamp() * 1000) for i in x: if "codebuild" not in i: print(i) await asyncio.run(export_to_s3(i,fromDate,toDate,destination_bucket)) async def export_to_s3(logGroupName1, fromTime1, to1, destination1): export_task = client.create_export_task( taskName = 'export_logs_to_s3', logGroupName = logGroupName1, fromTime = fromTime1, to = to1, destination = destination1, destinationPrefix = os.path.join(startTime.strftime('%Y{0}%m{0}%d').format(os.path.sep),logGroupName1) )
报错信息
接口返回错误:
{ "errorMessage": "Unable to marshal response: Object of type coroutine is not JSON serializable", "errorType": "Runtime.MarshalError", "requestId": "ca7a1113-b6a5-4ce9-a187-5867f1a91ea5", "stackTrace": [] }
函数运行日志:
Function Logs START RequestId: ca7a1113-b6a5-4ce9-a187-5867f1a91ea5 Version: $LATEST [ERROR] Runtime.MarshalError: Unable to marshal response: Object of type coroutine is not JSON serializable Traceback (most recent call last):/var/runtime/awslambdaric/bootstrap.py:405: RuntimeWarning: coroutine 'lambda_handler' was never awaited handle_event_request( RuntimeWarning: Enable tracemalloc to get the object allocation traceback END RequestId: ca7a1113-b6a5-4ce9-a187-5867f1a91ea5 REPORT RequestId: ca7a1113-b6a5-4ce9-a187-5867f1a91ea5 Duration: 6.86 ms Billed Duration: 7 ms Memory Size: 10240 MB Max Memory Used: 55 MB Init Duration: 291.92 ms Request ID ca7a1113-b6a5-4ce9-a187-5867f1a91ea5
问题根因
- 异步入口不兼容:AWS Lambda官方Python运行时不支持直接将异步协程(
async def)作为lambda_handler入口,运行时调用入口函数时拿到的是未执行的协程对象,无法将其序列化为JSON返回,直接触发序列化错误,日志中coroutine 'lambda_handler' was never awaited就是该问题的直接提示。 - 限流触发原因:CloudWatch Logs的
CreateExportTask接口有硬性账号级并发限制:同一区域同一时间仅允许1个导出任务处于运行/排队状态。原同步代码循环直接提交任务,前一个任务还未完成就提交下一个,直接触发LimitExceededException。 - 异步代码语法错误:
export_to_s3函数引用了外层作用域的client、startTime变量但未通过参数传递,就算协程能正常执行也会触发变量未定义错误。- boto3默认客户端是同步阻塞实现,套
async/await语法不会实现真正的异步执行,没有实际意义。 asyncio.run()方法会创建新的事件循环,不能在已经运行的事件循环中调用,在异步handler内循环调用该方法本身属于语法错误。
修复方案
不需要使用异步语法,直接用同步代码实现即可:每次提交导出任务后轮询任务状态,等当前任务执行完成(成功/失败)后再提交下一个任务,从根源上避开并发导出任务的数量限制。
修复后的完整可运行代码:
import boto3 import datetime import os import time def lambda_handler(event, context): client = boto3.client('logs') no_of_days = 7 destination_bucket = "abc" # 单次轮询任务状态的间隔时间,单位秒 poll_interval = 30 # 拉取所有日志组,支持分页 log_group_names = [] resp = client.describe_log_groups(limit=50) log_group_names.extend([lg['logGroupName'] for lg in resp['logGroups']]) while "nextToken" in resp: resp = client.describe_log_groups( nextToken=resp["nextToken"], limit=50 ) log_group_names.extend([lg['logGroupName'] for lg in resp['logGroups']]) # 计算导出时间范围 current_datetime = datetime.datetime.now() start_time = current_datetime - datetime.timedelta(days=no_of_days) from_ts = int(start_time.timestamp() * 1000) to_ts = int(current_datetime.timestamp() * 1000) export_results = [] for log_group_name in log_group_names: # 过滤不需要导出的日志组 if "codebuild" in log_group_name: continue print(f"开始处理日志组: {log_group_name}") # 构造S3存储路径前缀 dest_prefix = os.path.join( start_time.strftime('%Y{0}%m{0}%d').format(os.path.sep), log_group_name.lstrip('/') ) # 提交导出任务 task_resp = client.create_export_task( # 任务名使用唯一值,避免重名冲突 taskName = f'export_{log_group_name.replace("/","_")}_{int(time.time())}', logGroupName = log_group_name, fromTime = from_ts, to = to_ts, destination = destination_bucket, destinationPrefix = dest_prefix ) task_id = task_resp['taskId'] print(f"导出任务提交成功,任务ID: {task_id}") # 轮询任务状态,直到任务结束再处理下一个 while True: task_status_resp = client.describe_export_tasks(taskId=task_id) task_info = task_status_resp['exportTasks'][0] status_code = task_info['status']['code'] if status_code == 'COMPLETED': print(f"任务{task_id}执行完成") export_results.append({ "log_group": log_group_name, "task_id": task_id, "status": "SUCCESS" }) break if status_code in ['FAILED', 'CANCELLED']: err_msg = task_info['status'].get('message', '未知错误') print(f"任务{task_id}执行失败,原因: {err_msg}") export_results.append({ "log_group": log_group_name, "task_id": task_id, "status": "FAILED", "error": err_msg }) break # 任务仍在运行,等待后重试 time.sleep(poll_interval) return { "statusCode": 200, "export_count": len(export_results), "export_detail": export_results }
配置注意事项
- 权限配置:给Lambda执行角色附加以下权限:
logs:DescribeLogGroups、logs:CreateExportTask、logs:DescribeExportTasks,以及目标S3存储桶的对应写入权限。 - 超时配置:如果日志组数量多、单日志组数据量大,需要适当调大Lambda函数的超时时间(最长可配置15分钟);如果日志量过大,建议改成分批触发逻辑,每次调用仅导出1个日志组,避免超时。
- 导出任务本身是CloudWatch服务端异步执行的,不需要在代码层面使用异步语法,只要保证同一时间只有一个运行中的导出任务即可避开限流。
内容的提问来源于stack exchange,提问作者ASingh
相关产品推荐
相关产品推荐

