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

使用aiobotocore共享S3客户端异步下载文件遇协程重用错误

问题:aiobotocore重用客户端时报错"cannot reuse already awaited coroutine"

我想用aiobotocore异步上传下载S3文件,为避免每次请求新建客户端(速度较慢),决定重用客户端,编写了如下S3Utils类:

class S3Utils:

    def __init__(self):
        self._session = get_session()
        self._s3_client = self._session.create_client('s3', region_name=config.S3_REGION)

    async def upload_file(self, file_path):
        try:
            async with self._s3_client as s3_client:
                async with aiofiles.open(file_path, 'rb') as f:
                    await s3_client.put_object(Bucket=config.S3_BUCKET, Key=file_path, Body=(await f.read()))
        except Exception as e:
            logger.error('Error uploading file', exc_info=e)

    async def download_file(self, file_path):
        try:
            async with self._s3_client as s3_client:
                response = await s3_client.get_object(Bucket=config.S3_BUCKET, Key=file_path)
                content = await response['Body'].read()

            async with aiofiles.open(file_path, 'wb') as f:
                await f.write(content)

            return file_path
        except Exception as e:
            logger.error('Error downloading file', exc_info=e)
            return None

    async def download_all_files_in_directory(self, bucket, directory_path):
        try:
            files = []
            async with self._s3_client as s3_client:
                response = await s3_client.list_objects_v2(Bucket=bucket, Prefix=directory_path)
                if 'Contents' in response:
                    for file in response['Contents']:
                        file_path = file['Key']
                        if file_path.endswith(('jpg', 'png')):
                            files.append(file_path)

            results = await asyncio.gather(*[asyncio.create_task(self.download_file(file)) for file in files])
            
            print(f'results: {results}')
        except Exception as e:
            logger.error('Error downloading files', exc_info=e)

if __name__ == '__main__':
    asyncio.run(S3Utils().download_all_files_in_directory('bucket-name', 'folder-name'))

运行时出现以下异常:

2024-03-12T07:13:51.177643Z [error    ] Error downloading file         filename=s3_utils.py func_name=download_file lineno=38
Traceback (most recent call last):
  File "<path>\src\utils\s3_utils.py", line 29, in download_file
    async with self._s3_client as s3_client:
  File "<path>\venv\lib\site-packages\aiobotocore\session.py", line 25, in __aenter__
    self._client = await self._coro
RuntimeError: cannot reuse already awaited coroutine
results: [None, None, None, None, None, None, None, None, None, None, None, None]

将download_file的代码直接复制到download_all_files_in_directory方法内可正常运行,但代码冗余,想知道原写法为何报错及解决方法。


报错原因

self._session.create_client()返回的是协程对象,而非已初始化完成的客户端实例。第一次执行async with self._s3_client时,该协程会被自动await完成客户端初始化;但协程对象只能被await一次,后续再次通过async with尝试重用这个协程对象时,就会触发RuntimeError: cannot reuse already awaited coroutine错误。

解决方法

核心思路是提前完成客户端协程的await,存储已就绪的客户端实例,而非协程对象。以下两种方案都能实现客户端重用,同时避免报错:

方案1:延迟初始化客户端

在类中新增异步初始化方法,第一次调用业务方法时完成客户端的创建:

class S3Utils:
    def __init__(self):
        self._session = get_session()
        self._s3_client = None  # 先置空,后续延迟初始化

    async def _init_client(self):
        # 仅当客户端未初始化时执行创建
        if not self._s3_client:
            self._s3_client = await self._session.create_client('s3', region_name=config.S3_REGION)

    async def upload_file(self, file_path):
        try:
            await self._init_client()
            async with self._s3_client as s3_client:
                async with aiofiles.open(file_path, 'rb') as f:
                    await s3_client.put_object(Bucket=config.S3_BUCKET, Key=file_path, Body=(await f.read()))
        except Exception as e:
            logger.error('Error uploading file', exc_info=e)

    async def download_file(self, file_path):
        try:
            await self._init_client()
            async with self._s3_client as s3_client:
                response = await s3_client.get_object(Bucket=config.S3_BUCKET, Key=file_path)
                content = await response['Body'].read()

            async with aiofiles.open(file_path, 'wb') as f:
                await f.write(content)

            return file_path
        except Exception as e:
            logger.error('Error downloading file', exc_info=e)
            return None

    async def download_all_files_in_directory(self, bucket, directory_path):
        try:
            await self._init_client()
            files = []
            async with self._s3_client as s3_client:
                response = await s3_client.list_objects_v2(Bucket=bucket, Prefix=directory_path)
                if 'Contents' in response:
                    for file in response['Contents']:
                        file_path = file['Key']
                        if file_path.endswith(('jpg', 'png')):
                            files.append(file_path)

            results = await asyncio.gather(*[asyncio.create_task(self.download_file(file)) for file in files])
            
            print(f'results: {results}')
        except Exception as e:
            logger.error('Error downloading files', exc_info=e)

if __name__ == '__main__':
    asyncio.run(S3Utils().download_all_files_in_directory('bucket-name', 'folder-name'))

方案2:将类实现为异步上下文管理器

通过异步上下文管理器的__aenter__方法初始化客户端,__aexit__方法关闭客户端,确保生命周期管理规范:

class S3Utils:
    def __init__(self):
        self._session = get_session()
        self._s3_client = None

    async def __aenter__(self):
        # 进入上下文时初始化客户端
        self._s3_client = await self._session.create_client('s3', region_name=config.S3_REGION)
        return self

    async def __aexit__(self, exc_type, exc_val, exc_tb):
        # 退出上下文时关闭客户端
        await self._s3_client.close()

    async def upload_file(self, file_path):
        try:
            async with self._s3_client as s3_client:
                async with aiofiles.open(file_path, 'rb') as f:
                    await s3_client.put_object(Bucket=config.S3_BUCKET, Key=file_path, Body=(await f.read()))
        except Exception as e:
            logger.error('Error uploading file', exc_info=e)

    async def download_file(self, file_path):
        try:
            async with self._s3_client as s3_client:
                response = await s3_client.get_object(Bucket=config.S3_BUCKET, Key=file_path)
                content = await response['Body'].read()

            async with aiofiles.open(file_path, 'wb') as f:
                await f.write(content)

            return file_path
        except Exception as e:
            logger.error('Error downloading file', exc_info=e)
            return None

    async def download_all_files_in_directory(self, bucket, directory_path):
        try:
            files = []
            async with self._s3_client as s3_client:
                response = await s3_client.list_objects_v2(Bucket=bucket, Prefix=directory_path)
                if 'Contents' in response:
                    for file in response['Contents']:
                        file_path = file['Key']
                        if file_path.endswith(('jpg', 'png')):
                            files.append(file_path)

            results = await asyncio.gather(*[asyncio.create_task(self.download_file(file)) for file in files])
            
            print(f'results: {results}')
        except Exception as e:
            logger.error('Error downloading files', exc_info=e)

if __name__ == '__main__':
    async def main():
        async with S3Utils() as s3_utils:
            await s3_utils.download_all_files_in_directory('bucket-name', 'folder-name')
    asyncio.run(main())

两种方案都能解决协程重复await的问题,同时实现客户端重用,避免代码冗余。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 02:59:59