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

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
问题根因
  1. 异步入口不兼容:AWS Lambda官方Python运行时不支持直接将异步协程(async def)作为lambda_handler入口,运行时调用入口函数时拿到的是未执行的协程对象,无法将其序列化为JSON返回,直接触发序列化错误,日志中coroutine 'lambda_handler' was never awaited就是该问题的直接提示。
  2. 限流触发原因:CloudWatch Logs的CreateExportTask接口有硬性账号级并发限制:同一区域同一时间仅允许1个导出任务处于运行/排队状态。原同步代码循环直接提交任务,前一个任务还未完成就提交下一个,直接触发LimitExceededException。
  3. 异步代码语法错误:
    • 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 10:01:00