【问题标题】:how to start multiple jobs in python and communicate with the main job如何在python中启动多个作业并与主要作业进行通信
【发布时间】:2016-11-25 10:09:07
【问题描述】:

我是python多线程/多处理的新手,所以请多多包涵。 我想解决以下问题,在这方面我需要一些帮助/建议。 让我简单描述一下:

  1. 我想启动一个 python 脚本,它在 按顺序开始。

  2. 顺序部分结束后,我想开始一些工作 并行。

    • 假设我要启动四个并行作业。
    • 我还想在计算集群上使用“lsf”在其他一些机器上启动这些作业。我的初始脚本也在“lsf”上运行 机器。
    • 我在四台机器上开始的四个作业将依次执行两个逻辑步骤 A 和 B。
    • 当作业最初开始时,它们会从逻辑步骤 A 开始并完成它。
    • 在每个作业(4 个作业)完成步骤 A 之后;他们应该通知启动这些的第一个工作。也就是说,启动的主要作业就是等待这四个作业的确认。
    • 一旦主作业收到这四个作业的确认;它应该通知所有四个作业执行逻辑步骤 B。
    • 逻辑步骤 B 将在完成任务后自动终止作业。
    • 主要作业正在等待所有作业完成,稍后应继续执行后续部分。
  3. 一个示例场景是:

    • 在集群中的“lsf”机器上运行的 Python 脚本会在四个“lsf”机器上启动四个“tcl shell”。
    • 在每个 tcl shell 中,都有一个脚本用于执行逻辑步骤 A。
    • 一旦步骤 A 完成,他们应该以某种方式通知正在等待确认的 python 脚本。
    • 一旦收到四个人的确认,python 脚本就会通知他们执行逻辑步骤 B。
    • 逻辑步骤 B 也是一个脚本,来源于他们的 tcl shell;这个脚本最后也会关闭 tcl shell。
    • 同时,python 脚本正在等待所有四个作业完成。
    • 四个作业都完成后;它应该再次继续顺序部分并稍后完成。

这是我的问题:

  1. 我很困惑——我应该使用多线程/多处理。哪个更适合? 其实这两者有什么区别?我阅读了这些内容,但无法得出结论。

  2. 什么是 Python GIL?我还在任何时间点的某个地方读到,只有一个线程会执行。 我在这里需要一些解释。它给我的印象是我不能使用线程。

  3. 关于如何系统地以更 Python 的方式解决我的问题的任何建议。 我正在寻找一些逐步的口头解释和一些关于每一步的阅读指南。 一旦概念清楚,我想自己编写代码。

提前致谢。

【问题讨论】:

  • 看看celery。它应该在这里解决您的大部分/所有查询。

标签: python multithreading python-2.7 multiprocessing


【解决方案1】:

除了 roganjosh 的回答之外,我还会在 A 完成后添加一些信号以开始步骤 B:

import multiprocessing as mp
import time
import random
import sys

def func_A(process_number, queue, proceed):
    print "Process {} has started been created".format(process_number)

    print "Process {} has ended step A".format(process_number)
    sys.stdout.flush()
    queue.put((process_number, "done"))

    proceed.wait() #wait for the signal to do the second part
    print "Process {} has ended step B".format(process_number)
    sys.stdout.flush()

def multiproc_master():
    queue = mp.Queue()
    proceed = mp.Event()

    processes = [mp.Process(target=func_A, args=(x, queue)) for x in range(4)]
    for p in processes:
        p.start()

    #block = True waits until there is something available
    results = [queue.get(block=True) for p in processes]
    proceed.set() #set continue-flag
    for p in processes: #wait for all to finish (also in windows)
        p.join()
    return results

if __name__ == '__main__':
    split_jobs = multiproc_master()
    print split_jobs

【讨论】:

  • RaJa:你能解释一下为什么代码中有两个连接吗?
  • 嗯,这显然是一个错误。你只需要最后一个。我已经复制/粘贴了另一个解决方案并对其进行了一些修改。我已经编辑了我的解决方案并删除了错误的加入命令。
  • 不要在 Windows 上使用 .join()stackoverflow.com/questions/39896807/… 。充其量它会导致进程按顺序执行,最坏的情况是您最终会出现僵尸进程。
  • 您可能是对的,但我从未遇到过问题,而且我的流程显然是并行运行的。我正在使用 Windows 和 Python 2.7/3。
  • Roganjosh/RaJa:我明白你的提议。在 func_A 中,您计算​​随机选择的数字的平方和。这是您尝试过的一个简单的 python 函数。就我而言,我想在远程机器上打开“tclsh”外壳并获取 tcl 脚本。此脚本需要一些时间才能运行,例如 1 小时。我应该使用 subprocess 来启动 tclshell 吗?并使用管道从标准输出和标准错误中读取?如何识别步骤A已经完成?我是否应该在标准输出中查找一些特定的词来确定步骤 A 已完成?
【解决方案2】:

1) 根据您在问题中列出的选项,在这种情况下,您可能应该使用 multiprocessing 来利用多个 CPU 内核并并行计算事物。

2)从第1点更进一步:全局解释器锁(GIL)意味着在任何时候只有一个线程可以实际执行代码。
这里经常弹出的multithreading 的一个简单示例是提示用户输入例如数学问题的答案。在后台,他们想要一个计时器以一秒的间隔保持递增,以记录该人响应所需的时间。如果没有多线程,程序会在等待用户输入时阻塞,并且计数器不会增加。在这种情况下,您可以让计数器和输入提示在不同的线程上运行,以便它们看起来同时运行。
实际上,两个线程共享相同的 CPU 资源,并且不断地前后传递一个对象(GIL)以授予它们对 CPU 的单独访问权限。如果您想正确地并行处理事情,这是没有希望的。 (注意:实际上,您只需记录提示前后的时间并计算差异,而不是打扰线程。)

3) 我用multiprocessing 做了一个非常简单的例子。在这种情况下,我生成了 4 个进程来计算随机选择的范围的平方和。这些进程没有共享的 GIL,因此与multithreading 不同的是独立执行。在此示例中,您可以看到所有进程的开始和结束时间略有不同,但我们可以将进程的结果聚合到单个 queue 对象中。在继续之前,父进程将等待所有 4 个子进程返回它们的计算。然后,您可以重复 func_B 的代码(未包含在代码中)。

import multiprocessing as mp
import time
import random
import sys

def func_A(process_number, queue):
    start = time.time()
    print "Process {} has started at {}".format(process_number, start)
    sys.stdout.flush()
    my_calc = sum([x**2 for x in xrange(random.randint(1000000, 3000000))])

    end = time.time()
    print "Process {} has ended at {}".format(process_number, end)
    sys.stdout.flush()
    queue.put((process_number, my_calc))

def multiproc_master():
    queue = mp.Queue()

    processes = [mp.Process(target=func_A, args=(x, queue)) for x in xrange(4)]
    for p in processes:
        p.start()

    # Unhash the below if you run on Linux (Windows and Linux treat multiprocessing
    # differently as Windows lacks os.fork())
    #for p in processes:
    #    p.join()

    results = [queue.get() for p in processes]
    return results

if __name__ == '__main__':
    split_jobs = multiproc_master()
    print split_jobs

【讨论】:

  • Roganjosh:我明白你的提议。在 func_A 中,您计算​​随机选择的数字的平方和。这是您尝试过的一个简单的 python 函数。就我而言,我想在远程机器上打开“tclsh”外壳并获取 tcl 脚本。此脚本需要一些时间才能运行,例如 1 小时。我应该使用 subprocess 来启动 tclshell 吗?并使用管道从标准输出和标准错误中读取?如何识别步骤A已经完成?我是否应该在标准输出中查找一些特定的词来确定步骤 A 已完成?
  • 这里函数的复杂性我认为不相关。敲几个 0,它的行为仍然相同,即处理时间不是问题。此代码隐式标识所有子 func_A 调用的结束,因为 results 在所有进程返回之前不会计算。我认为除了我在这里的内容之外,您可能还想查看阻塞/非阻塞子进程调用stackoverflow.com/questions/21936597/…。您可以从多进程中生成阻塞进程,或者非阻塞进程并在它们全部完成时设置一些标志
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2013-03-15
  • 1970-01-01
  • 2020-11-13
  • 2012-08-19
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多