Python:类似于`map`的东西,适用于线程

Ram*_*hum 26 python parallel-processing multithreading map-function

我确信在标准库中有这样的东西,但似乎我错了.

我有一堆我想urlopen并行的网址.我想要内置map函数,除了工作由一堆线程并行完成.

是否有一个很好的模块可以做到这一点?

Sco*_*son 42

多处理中有一种map方法.Pool.这会做多个过程.

如果多个进程不是你的菜,你可以使用使用线程的multiprocessing.dummy.

import urllib
import multiprocessing.dummy

p = multiprocessing.dummy.Pool(5)
def f(post):
    return urllib.urlopen('http://stackoverflow.com/questions/%u' % post)

print p.map(f, range(3329361, 3329361 + 5))
Run Code Online (Sandbox Code Playgroud)

  • 惊人的。值得注意的是,它下面实际上是`from multiprocessing.pool import ThreadPool` (3认同)

Ram*_*hum 11

有人建议我使用这个futures包.我试了一下它似乎工作.

http://pypi.python.org/pypi/futures

这是一个例子:

"Download many URLs in parallel."

import functools
import urllib.request
import futures

URLS = ['http://www.foxnews.com/',
        'http://www.cnn.com/',
        'http://europe.wsj.com/',
        'http://www.bbc.co.uk/',
        'http://some-made-up-domain.com/']

def load_url(url, timeout):
    return urllib.request.urlopen(url, timeout=timeout).read()

with futures.ThreadPoolExecutor(50) as executor:
   future_list = executor.run_to_futures(
           [functools.partial(load_url, url, 30) for url in URLS])
Run Code Online (Sandbox Code Playgroud)


Rug*_*nar 5

这是我对线程映射的实现:

from threading import Thread
from queue import Queue

def thread_map(f, iterable, pool=None):
    """
    Just like [f(x) for x in iterable] but each f(x) in a separate thread.
    :param f: f
    :param iterable: iterable
    :param pool: thread pool, infinite by default
    :return: list if results
    """
    res = {}
    if pool is None:
        def target(arg, num):
            try:
                res[num] = f(arg)
            except:
                res[num] = sys.exc_info()

        threads = [Thread(target=target, args=[arg, i]) for i, arg in enumerate(iterable)]
    else:
        class WorkerThread(Thread):
            def run(self):
                while True:
                    try:
                        num, arg = queue.get(block=False)
                        try:
                            res[num] = f(arg)
                        except:
                            res[num] = sys.exc_info()
                    except Empty:
                        break

        queue = Queue()
        for i, arg in enumerate(iterable):
            queue.put((i, arg))

        threads = [WorkerThread() for _ in range(pool)]

    [t.start() for t in threads]
    [t.join() for t in threads]
    return [res[i] for i in range(len(res))]
Run Code Online (Sandbox Code Playgroud)