【问题标题】:How to implement Producer Consumer with multiprocessing?如何通过多处理实现生产者消费者?
【发布时间】:2018-04-23 04:00:45
【问题描述】:

我有一个程序需要从某个来源下载文件并上传。但我需要确保下载位置最多有 10 个文件。有没有办法使用 Managers() 吗?

这听起来像是典型的生产者 - 消费者问题。下面是我的程序。

下面是我的实现

from multiprocessing import Process, Queue, Lock
import requests
import json
import shutil
import os
import time
import random
import warnings
warnings.filterwarnings("ignore")

sha_list = [line.strip() for line in open("ShaList")]


def save_file_from_sofa(sha1):
    r = requests.get("https://DOWNLOAD_URL/{}".format(sha1), verify=False, stream=True)
    with open(sha1, 'wb') as handle:
        shutil.copyfileobj(r.raw, handle)


def mock_upload():
    time.sleep(random.randint(10,16))


def producer(queue, lock):
    with lock:
        print("Starting Producer {}".format(os.getpid()))

    while sha_list:
        if not queue.full():
            sha1 = sha_list.pop()
            save_file_from_sofa(sha1)
            queue.put(sha1)


def consumer(queue, lock):
    with lock:
        print("Starting Consumer {}".format(os.getpid()))

    while True:
        sha1 = queue.get()
        mock_upload()
        with lock:
            print("{} GOT {}".format(os.getpid(), sha1))

if __name__ == "__main__":
    queue = Queue(5)
    lock = Lock()

    producers = [Process(target=producer, args=(queue, lock)) for _ in range(2)]
    consumers = []

    for _ in range(3):
        p = Process(target=consumer, args=(queue, lock))
        p.daemon = True #Do not forget to set it to true
        consumers.append(p)

    for p in producers:
        p.start()
    for c in consumers:
        c.start()

    for p in producers:
        p.join()

    print("DONE")

但它并没有达到预期的效果,正如您从下面的输出中看到的那样

启动生产者 623

启动生产者 624

启动消费者 626

启动消费者 625

启动消费者 627

626 GOT 4ff551490d6b2eec7c6c0470f4b092fdc34cd521

625 GOT 83a53a3400fc83f2b02135ba0cc6c8625ecc7dc4

627 GOT 4ff551490d6b2eec7c6c0470f4b092fdc34cd521

626 GOT 83a53a3400fc83f2b02135ba0cc6c8625ecc7dc4

625 GOT 4e7132301ce9d61445db07910ff90a64474e6a88

626 GOT 0efbd413d733b3903e6dee777ace5ef47a2ec144

627 GOT 4e7132301ce9d61445db07910ff90a64474e6a88

625 GOT 0efbd413d733b3903e6dee777ace5ef47a2ec144

626 GOT 0a3fc4bdd56fa2bf52f5f43277f3b4ee0f040937

625 GOT eb9c07329a8b5cb66e47f0dd8e56894707a84d94

627 GOT 0a3fc4bdd56fa2bf52f5f43277f3b4ee0f040937

626 GOT eb9c07329a8b5cb66e47f0dd8e56894707a84d94

完成

如您所见,消费者多次获取相同的 SHA1。所以,我需要一个程序来确保生产者放入队列中的所有 SHA1 仅被 1 个消费者拾取。

P.S 我也想过使用 pool 让它工作。对于生产者来说,它可以正常工作,因为我已经将 SHA1 列表放入队列,但是对于消费者,我将如何使用任何列表来确保消费者实际上正在停止。

【问题讨论】:

    标签: python multiprocessing python-multiprocessing


    【解决方案1】:

    只需使用来自multiprocessing.Pool 或concurrent.futures 的池。池允许您设置要同时运行的工作人员数量。这意味着您最多可以同时下载max_workers 个文件。

    由于下载/上传是连续的(在下载完成之前您无法开始上传),因此在两个单独的线程/进程中运行它们没有任何价值。只需将这两个操作加入一个作业单元中,然后同时运行多个作业。

    此外,只要您只需要下载/上传文件(IO 绑定操作),您最好使用线程而不是进程,因为它们更轻量级。

    from concurrent.futures import ThreadPoolExecutor
    
    list_of_sha1s = ['foobar', 'foobaz']
    
    def worker(sha1):
        path = save_file_from_sofa(sha1)
        upload_file(path)
    
        return sha1
    
    with ThreadPoolExecutor(max_workers=10) as pool:
        for sha1 in pool.map(worker, list_of_sha1s):
            print("Done SHA1: %s" % sha1)
    

    【讨论】:

    • 不是顺序操作。这就是我使用队列的原因。下载速度比上传快3倍。所以我希望有 2 个下载进程和 3,4 个上传进程。下载文件的 SHA1 被转储到队列中,消费者从那里选择它。它继续......@noxdafox
    • 是的,不能解决生产者-消费者问题,即你有一个单独的进程生产到一个消费。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2012-11-14
    • 1970-01-01
    • 1970-01-01
    • 2018-08-26
    • 1970-01-01
    • 1970-01-01
    • 2010-10-29
    相关资源
    最近更新 更多