如何使用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
相关产品推荐
相关产品推荐

