【问题标题】:Python 3 concurrent.futures: How to add back failed futures to ThreadPoolExecutor?Python 3 concurrent.futures:如何将失败的期货添加回 ThreadPoolExecutor?
【发布时间】:2014-10-14 04:10:19
【问题描述】:

我有一个要通过 concurrent.futures 的 ThreadPoolExecutor 下载的 url 列表,但可能有一些超时 url,我想在所有第一次尝试结束后重新下载它们。我不知道该怎么做,这是我的尝试,但因无休止的打印“time_out_again”而失败:

import concurrent.futures

def player_url(url):
    # here. if timeout, return 1. otherwise do I/O and return 0.
    ...

urls = [...]
time_out_futures = [] #list to accumulate timeout urls
with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor:
    future_to_url = (executor.submit(player_url, url) for url in urls)
    for future in concurrent.futures.as_completed(future_to_url):
        if future.result() == 1:
            time_out_futures.append(future)

# here is what I try to deal with all the timeout urls       
while time_out_futures:
    future = time_out_futures.pop()
    if future.result() == 1:
        print('time_out_again')
        time_out_futures.insert(0,future)   # add back to the list

那么,有什么办法可以解决这个问题吗?

【问题讨论】:

    标签: python concurrent.futures


    【解决方案1】:

    Future 对象只能使用一次。 Future 本身对返回结果的函数一无所知 - ThreadPoolExecutor 对象负责创建 Future、返回它并在后台运行函数:

    def submit(self, fn, *args, **kwargs):
        with self._shutdown_lock:
            if self._shutdown:
                raise RuntimeError('cannot schedule new futures after shutdown')
    
            f = _base.Future()
            w = _WorkItem(f, fn, args, kwargs)
    
            self._work_queue.put(w)
            self._adjust_thread_count()
            return f
    
    class _WorkItem(object):
        def __init__(self, future, fn, args, kwargs):
            self.future = future
            self.fn = fn
            self.args = args
            self.kwargs = kwargs
    
        def run(self):
            if not self.future.set_running_or_notify_cancel():
                return
    
            try:
                result = self.fn(*self.args, **self.kwargs)  # sefl.fn is play_url in your case
            except BaseException as e:
                self.future.set_exception(e)
            else:
                self.future.set_result(result)  # The result is set on the Future
    

    如您所见,函数完成后,结果将设置在Future 对象上。因为Future 对象实际上对提供结果的函数一无所知,所以无法尝试使用Future 对象重新运行该函数。你所能做的就是在超时发生时将url1一起返回,然后将submit的url重新发送到ThreadPoolExecutor

    def player_url(url):
        # here. if timeout, return 1. otherwise do I/O and return 0.
        ...
        if timeout:
            return (1, url)
        else:
            return (0, url)
    
    urls = [...]
    with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor:
        while urls:
            future_to_url = executor.map(player_url, urls)
            urls = []  # Clear urls list, we'll re-add any timed out operations.
            for future in future_to_url:
                if future.result()[0] == 1:
                    urls.append(future.result()[1]) # stick url into list
    

    【讨论】:

    • 只是一个小问题,如果你使用executor.map(player_url, urls),稍后会提出AttributeError: 'tuple' object has no attribute '_condition'。我用原来的方式,效果很好。
    • @xiang 啊,对不起。我已经纠正了这个错误。 as_completed 不应与 map 返回的迭代器一起使用。使用你原来的方法也很好。这种方式只是更简洁一点。
    猜你喜欢
    • 2022-01-10
    • 1970-01-01
    • 1970-01-01
    • 2013-08-26
    • 2018-05-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-01-10
    相关资源
    最近更新 更多