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

如何使用aioboto3与asyncio异步下载AWS S3文件(Python)

报错根因

aioboto3 和同步版 boto3 的资源创建逻辑完全不同:调用 session.resource() / session.client() 拿到的不是可直接使用的资源/客户端实例,而是一个ResourceCreatorContext异步上下文管理器对象,必须通过async with语法进入上下文后,才能拿到实际绑定了连接的可用对象,直接在上下文外部调用.Queue()方法就会触发你遇到的 'ResourceCreatorContext' object has no attribute 'Queue' 报错。

完整可运行异步改造代码

首先安装依赖:

pip install aioboto3

改造后的完整代码:

import json
import io
import asyncio
import gzip
import logging
from logging.handlers import RotatingFileHandler
import aioboto3

AWS_KEY = "**"
AWS_SECRET = "**"
QUEUE_URL = "***"
OUTPUT_PATH = "./test"
VISIBILITY_TIMEOUT = 30  # 建议设置为单条消息最大处理时长的1.5倍,避免消息处理中途被重投
REGION_NAME = "region"
SLEEP_TIME = 1
MAX_CONCURRENT_FILE = 5  # 控制S3文件并发下载数,避免触发AWS接口限流

sem = asyncio.Semaphore(MAX_CONCURRENT_FILE)

async def handle_response(msg, path):
    """业务处理逻辑,如有IO操作需替换为异步实现"""
    print(f'message from {path}: {msg}')

async def download_single_file(s3_client, bucket, s3_path):
    async with sem:
        with io.BytesIO() as f:
            # 异步下载S3 gzip文件
            await s3_client.download_fileobj(bucket, s3_path, f)
            f.seek(0)
            # 逐行解压处理,gzip解压为CPU操作,大文件场景可丢入线程池执行避免阻塞事件循环
            for line in gzip.GzipFile(fileobj=f):
                await handle_response(line.decode('UTF-8'), s3_path)

async def consume():
    session = aioboto3.Session()
    # 所有客户端/资源必须在async with上下文内初始化
    async with session.resource(
        'sqs',
        region_name=REGION_NAME,
        aws_access_key_id=AWS_KEY,
        aws_secret_access_key=AWS_SECRET
    ) as sqs, session.client(
        's3',
        region_name=REGION_NAME,
        aws_access_key_id=AWS_KEY,
        aws_secret_access_key=AWS_SECRET
    ) as s3:
        queue = await sqs.Queue(url=QUEUE_URL)
        logging.info("SQS consumer started, listening on queue: %s", QUEUE_URL)
        while True:
            # 异步拉取消息,MaxNumberOfMessages可根据消费能力设置为1-10
            messages = await queue.receive_messages(
                VisibilityTimeout=VISIBILITY_TIMEOUT,
                MaxNumberOfMessages=10
            )
            if not messages:
                await asyncio.sleep(SLEEP_TIME)
                continue
            # 并发处理拉取到的消息
            tasks = []
            for msg in messages:
                async def process_single_message(msg):
                    try:
                        body = json.loads(await msg.body)
                        # 并发下载当前消息关联的所有S3文件
                        file_tasks = [
                            download_single_file(s3, body['bucket'], s3_file['path'])
                            for s3_file in body['files']
                        ]
                        await asyncio.gather(*file_tasks)
                        # 处理完成后删除消息
                        await msg.delete()
                    except Exception as e:
                        logging.error("Process message failed: %s", str(e), exc_info=True)
                tasks.append(process_single_message(msg))
            await asyncio.gather(*tasks)

if __name__ == '__main__':
    # 日志初始化
    logging.basicConfig(level=logging.INFO, format="%(asctime)s %(name)s %(levelname)s %(message)s")
    logger = logging.getLogger("Consumer")
    RFH = RotatingFileHandler("test.log", maxBytes=20971520, backupCount=5)
    F_FORMAT = logging.Formatter('%(asctime)s %(name)s %(levelname)s %(message)s')
    RFH.setLevel(logging.INFO)
    RFH.setFormatter(F_FORMAT)
    logger.addHandler(RFH)
    # 启动异步事件循环
    asyncio.run(consume())
关键改造注意点
  • 所有aioboto3返回的资源、客户端对象必须在async with块内使用,禁止在上下文外部调用对象方法
  • 所有异步IO操作(消息接收、消息删除、S3下载、等待休眠)必须加await关键字,禁止在异步代码中使用同步阻塞方法(比如time.sleep、同步boto3调用),否则会卡住整个事件循环,异步优化完全失效
  • 增加并发控制信号量asyncio.Semaphore,避免无限制并发请求触发AWS的S3/SQS接口限流(默认单账号S3 QPS为3500,SQS为300,可根据账号配额调整并发数)
  • 原代码中VISIBILITY_TIMEOUT设置为10秒过短,如果单条消息关联的S3文件较大、处理时长超过10秒,消息会被SQS重新投递给其他消费者,导致重复消费,建议根据实际单条消息最大处理时长调整为1.5-2倍
  • 如果gzip文件单个体积超过100M,逐行解压的CPU操作会阻塞事件循环,可以把解压逻辑用asyncio.to_thread()丢到默认线程池执行,避免影响其他并发任务调度

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 22:01:13