ImageBatchGenerator类并发运行时出现重复加载URL问题
问题分析
你的ImageBatchGenerator出现重复加载同一URL的核心原因是没有跟踪已处理完成的Future:
- 在
__iter__方法中,你一次性提交了所有URL的请求任务,生成了所有futures。 - 每次调用
__next__时,你都会通过as_completed(self.futures)遍历所有Future——包括上一个批次已经处理过的那些。这些已完成的Future会被重复取出结果,自然就出现了重复的URL和图像。
修复方案
我们需要维护一个未处理的Future集合,每次处理完一个Future就从集合中移除,确保每个任务只被处理一次。下面是修正后的完整代码:
import requests from io import BytesIO from typing import Union, Set from concurrent.futures import ThreadPoolExecutor, Future, as_completed from PIL import Image def fetch_image(url: str) -> tuple[str, Image.Image]: """Fetch image from url Parameters ---------- url : str url of the image Returns ------- tuple[str, Image.Image] tuple (url, image) where image is PIL image object and url is the url of the image """ try: response = requests.get(url) response.raise_for_status() return url, Image.open(BytesIO(response.content)) except requests.exceptions.HTTPError as errh: print(f"HTTP Error: {errh}") except requests.exceptions.ConnectionError as errc: print(f"Error Connecting: {errc}") except requests.exceptions.Timeout as errt: print(f"Timeout Error: {errt}") except requests.exceptions.RequestException as err: print(f"Something Else: {err}") return None class ImageBatchGenerator: """ A generator class that get's as arguments a list of URLs and batch size and generates batches of PIL images that are obtained through GET requests to the URLs. Parameters ---------- urls : list[str] List of URLs to fetch images from batch_size : int The size of the batches to be generated """ def __init__(self, urls: list[str], batch_size: int=32) -> None: self.urls = urls self.batch_size = batch_size self.executor = ThreadPoolExecutor() def __len__(self) -> int: return (len(self.urls) + self.batch_size - 1) // self.batch_size def __iter__(self) -> ImageBatchGenerator: # 提交所有任务,并将futures存储为一个可修改的集合 self.pending_futures: Set[Future] = { self.executor.submit(fetch_image, url) for url in self.urls } return self def __next__(self) -> dict[str, Union[str, Image.Image]]: images = [] urls = [] # 遍历已完成的任务,同时从pending集合中移除 for future in as_completed(self.pending_futures): self.pending_futures.remove(future) result = future.result() if result is not None: url, image = result images.append(image) urls.append(url) # 凑够batch_size就停止当前批次的收集 if len(images) == self.batch_size: break if len(images) == 0: self.executor.shutdown() raise StopIteration return {"images": images, "urls": urls}
关键修改点
- 维护待处理Future集合:
- 在
__iter__中,将futures存储为一个Set[Future]而不是列表,方便快速移除已处理的任务。
- 在
- 每次处理后移除已完成的Future:
- 在
__next__中,处理完一个future后立即从pending_futures集合中移除它,确保不会被后续的批次重复处理。
- 在
- 正确处理
fetch_image的返回值:- 原代码中直接解包
url, image = future.result(),但如果fetch_image抛出异常后返回None,这会导致解包错误。现在先判断result是否为None,再进行解包,避免崩溃。
- 原代码中直接解包
这样修改后,每个URL对应的任务只会被处理一次,就不会再出现重复加载的问题了。
内容的提问来源于stack exchange,提问作者Eduardo Pacheco
相关产品推荐
相关产品推荐

