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

