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

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}
关键修改点
  1. 维护待处理Future集合:
    • 在__iter__中,将futures存储为一个Set[Future]而不是列表,方便快速移除已处理的任务。
  2. 每次处理后移除已完成的Future:
    • 在__next__中,处理完一个future后立即从pending_futures集合中移除它,确保不会被后续的批次重复处理。
  3. 正确处理fetch_image的返回值:
    • 原代码中直接解包url, image = future.result(),但如果fetch_image抛出异常后返回None,这会导致解包错误。现在先判断result是否为None,再进行解包,避免崩溃。

这样修改后,每个URL对应的任务只会被处理一次,就不会再出现重复加载的问题了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 16:50:16