Python中AWS Boto分页多线程失效问题排查及可行性咨询
问题:S3分页器与多线程结合失效,仅单线程运行
已知信息:
- 分页响应类型:
<class 'botocore.paginate.PageIterator'> page.get('Contents')返回包含对象名称和大小的字典列表- S3的
list_objects_v2接口单次最多返回1000个对象,因此必须通过分页器遍历所有对象
尝试将分页迭代器分配给线程分摊负载,但实际整个迭代过程仍为单线程运行。以下是两种尝试的实现代码:
第一种实现代码
def list_s3_files_using_paginator_multithreading(bucket_name, prefix = None): import boto3 from threading import Thread from queue import Queue """ This functions list all files in s3 using paginator. Paginator is useful when you have 1000s of files in S3. S3 list_objects_v2 can list at max 1000 files in one go. :return: None """ def safe_dedicated_writing_thread(filepath, queue): with open(filepath, 'w') as f: while True: # Retrieve s3 ls object name from the queue line = queue.get() # check if we have reached end of s3 ls dump # then we close the file if line is None: break # Write to file f.write(line) # In Python, files are automatically flushed # while closing them but, a programmer can flush # a file before closing it f.flush() # Mark the unit of work complete queue.task_done() # mark the exit signal as processed, after the file was closed queue.task_done() # Create a shared queue queue = Queue() # Path of the shared file filepath = 's3_ls_dump.txt' # Create And Start the file writer Thread writer_thread = Thread(target = safe_dedicated_writing_thread, args = (filepath, queue), daemon = True) writer_thread.start() # Creating a Paginator to list 1000 object per response s3_client = boto3.client("s3") paginator = s3_client.get_paginator("list_objects_v2") response = paginator.paginate(Bucket=bucket_name, Prefix = prefix, PaginationConfig={"PageSize": 1000}) def list_object(page): files = page.get('Contents') for file in files: queue.put(file['Key']) # Configure worker Thread with concurrent.futures.ThreadPoolExecutor() as executor: executor.map(list_object, response) # Signal the file writer thread that we are done queue.put(None) # Wait for all the queue to be processed queue.join()
第二种实现代码
def list_object(page, queue): files = page.get('Contents') for file in files: queue.put(file['Key']) threads = [Thread(target = list_object, args = (page, queue)) for page in response] for thread in threads: thread.start() for thread in threads: thread.join()
请问这种场景下是否可以实现真正的多线程?
解答
为什么之前的代码是单线程运行
botocore.paginate.PageIterator的迭代过程是单线程阻塞式的:无论是用executor.map遍历,还是生成线程列表时遍历,主线程都会先逐个从S3拉取分页数据(串行发起请求),拿到page之后才会交给线程处理。也就是说,分页请求本身是串行的,线程只是处理已经拉取到本地的page数据,并没有并发发起分页请求,所以整体看起来还是单线程运行。
这种场景可以实现真正的多线程
想要真正分摊负载,需要让多个线程并发发起S3分页请求,而不是串行拉取page之后再处理。核心思路是手动维护分页标记(ContinuationToken),让多个线程各自携带token发起请求,避免串行拉取。
正确的实现示例
import boto3 from threading import Thread from queue import Queue import concurrent.futures def list_s3_files_concurrent_pagination(bucket_name, prefix=None): # 负责写入文件的线程 def safe_writer_thread(filepath, queue): with open(filepath, 'w') as f: while True: line = queue.get() if line is None: break f.write(f"{line}\n") f.flush() queue.task_done() queue.task_done() # 初始化数据队列和写入线程 data_queue = Queue() output_file = 's3_ls_dump.txt' writer_thread = Thread(target=safe_writer_thread, args=(output_file, data_queue), daemon=True) writer_thread.start() s3_client = boto3.client("s3") # 用于存放分页标记的队列 token_queue = Queue() # 先拉取第一页数据,初始化分页标记队列 first_page = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=prefix, MaxKeys=1000) if 'Contents' in first_page: for obj in first_page['Contents']: data_queue.put(obj['Key']) # 如果有下一页,把标记放入队列 if 'NextContinuationToken' in first_page: token_queue.put(first_page['NextContinuationToken']) # 线程处理分页标记的逻辑:拉取对应页数据,并把下一个标记放回队列 def process_pagination_token(token): while token is not None: response = s3_client.list_objects_v2( Bucket=bucket_name, Prefix=prefix, MaxKeys=1000, ContinuationToken=token ) # 处理当前页的对象 if 'Contents' in response: for obj in response['Contents']: data_queue.put(obj['Key']) # 获取下一个分页标记,没有则设为None token = response.get('NextContinuationToken') if token is not None: token_queue.put(token) # 启动线程池,并发处理分页请求 with concurrent.futures.ThreadPoolExecutor(max_workers=4) as executor: # 持续从标记队列取任务,直到队列为空 while not token_queue.empty(): current_token = token_queue.get() executor.submit(process_pagination_token, current_token) token_queue.task_done() # 等待所有标记处理完成 token_queue.join() # 通知写入线程结束 data_queue.put(None) data_queue.join() # 调用示例 # list_s3_files_concurrent_pagination('你的桶名称', '你的前缀/')
代码说明
- 主线程先拉取第一页数据,并将后续分页标记放入
token_queue - 线程池中的多个线程从
token_queue取标记,各自发起S3分页请求,实现并发拉取 - 每个线程处理完当前页后,若有下一页标记,会将其放回队列,供其他线程继续处理
- 写入线程负责将所有对象Key写入文件,与拉取过程解耦
内容的提问来源于stack exchange,提问作者Himanshu Rajput
相关产品推荐
相关产品推荐

