您展示的recursive 生成器实际上并不是递归的,这会导致系统递归限制出现问题。
要了解为什么您需要注意recursive 生成器的代码何时运行。与普通函数不同,仅调用 recursive(0) 不会导致它立即运行其代码并进行额外的递归调用。相反,调用recursive(0) 会立即返回一个生成器对象。只有当你 send() 到生成器时,代码才会运行,只有在你第二次 send() 到它之后,它才会启动另一个调用。
让我们在代码运行时检查调用堆栈。在顶层,我们运行Task(recursive(0)).step()。依次执行三件事:
-
recursive(0) 此调用立即返回一个生成器对象。
-
Task(_) Task 对象被创建,它的__init__ 方法存储了对第一步创建的生成器对象的引用。
-
_.step() 任务上的一个方法被调用。这是行动真正开始的地方!让我们看看调用内部发生了什么:
-
fut = self._gen.send(value) 在这里,我们实际上是通过发送一个值来启动生成器运行。让我们更深入地查看生成器代码的运行情况:
-
yield pool.submit(time.sleep, 0.001) 这安排了在另一个线程中完成的事情。我们不会等待它发生。相反,我们会得到一个Future,我们可以使用它在完成时收到通知。我们将 future 立即返回到之前的代码级别。
-
fut.add_done_callback(self._wakeup) 在这里,我们要求在未来准备好时调用我们的 _wakeup() 方法。这总是立即返回!
-
step 方法现在结束。没错,我们已经完成了(暂时)!这对于您问题的第二部分很重要,稍后我将详细讨论。
我们的调用结束了,所以如果我们以交互方式运行,控制流将返回到 REPL。如果我们作为脚本运行,解释器将到达脚本的末尾并开始关闭(我将在下面详细讨论)。但是,由线程池控制的其他线程仍在运行,并且在某些时候,其中一个会做一些我们关心的事情!让我们看看那是什么。
-
当计划函数 (time.sleep) 完成运行后,它正在运行的线程将调用我们在 Future 对象上设置的回调。也就是说,它将在我们之前创建的Task 对象上调用Task._wakup()(我们在顶层不再引用它,但Future 保留了一个引用,因此它仍然存在)。我们来看看方法:
-
result = fut.result() 存储延迟调用的结果。在这种情况下,这无关紧要,因为我们从不查看结果(无论如何都是None)。
-
self.step(result) 再接再厉!现在我们回到我们关心的代码。让我们看看它这次做了什么:
-
fut = self._gen.send(value) 再次发送给生成器,所以它接管了。它已经产生了一次,所以这次我们在 yield 之后开始:
-
print("Tick :", n) 这个很简单。
-
Task(recursive(n+1)).step() 这就是事情变得有趣的地方。这条线就像我们开始的那样。所以,和以前一样,这将运行我上面列出的逻辑 1-4(包括它们的子步骤)。但是,当step() 方法返回时,它不会返回到 REPL 或结束脚本,而是返回到这里。
-
recursive() 生成器(原来的生成器,不是我们刚刚创建的新生成器)已经结束。因此,就像任何到达其代码末尾的生成器一样,它会引发StopIteration。
-
StopIteration 被 try/except 块捕获并忽略,step() 方法结束。
-
_wakup() 方法也结束,所以回调完成。
- 最终,在前面的回调中创建的
Task 的回调也将被调用。所以我们返回并重复第 5 步,一遍又一遍,永远(如果我们以交互方式运行)。
上面的调用堆栈解释了为什么交互式案例永远打印出来。主线程返回到 REPL(如果你能看到其他线程的输出,你可以用它做其他事情)。但是在池中,每个线程从自己的作业的回调中调度另一个作业。当下一个作业完成时,它的回调会安排另一个作业,依此类推。
那么,当您将代码作为脚本运行时,为什么只得到 8 个打印输出?上面的第 4 步暗示了答案。当以非交互方式运行时,主线程在第一次调用Task.step 返回后运行结束脚本。这会提示解释器尝试关闭。
concurrent.futures.thread 模块(其中定义了ThreadPoolExecutor)有一些奇特的逻辑,当程序关闭而执行程序仍处于活动状态时,它会尝试很好地清理。它应该停止任何空闲线程,并在当前作业完成时向仍在运行的任何线程发出信号以停止。
该清理逻辑的确切实现以一种非常奇怪的方式与我们的代码交互(可能有问题,也可能没有问题)。效果是第一个线程不断地给自己更多的工作要做,而产生的其他工作线程在产生后立即退出。当 executor 启动了尽可能多的线程(在我们的例子中是 8 个)时,第一个 worker 最终退出。
据我所知,这是事件的顺序。
- 我们(间接)导入
concurrent.futures.thread 模块,该模块使用atexit 告诉解释器在解释器关闭之前运行一个名为_python_exit 的函数。
- 我们创建了一个最大线程数为 8 的
ThreadPoolExecutor。它不会立即产生其工作线程,而是会在每次调度作业时创建一个,直到它拥有全部 8 个线程。
- 我们安排我们的第一个作业(在上一个列表中第 3 步的深层嵌套部分)。
- 执行程序将作业添加到其内部队列,然后注意到它没有最大数量的工作线程并启动一个新线程。
- 新线程将作业从队列中弹出并开始运行。但是,
sleep 调用比其余步骤花费的时间要长得多,因此线程会在这里卡住一段时间。
- 主线程完成(已到达上一个列表中的第 4 步)。
-
_python_exit 函数被解释器调用,因为解释器想要关闭。该函数在模块中设置一个全局_shutdown 变量,并将None 发送到执行器的内部队列(它为每个线程发送一个None,但目前只创建了一个线程,所以它只是发送一个None)。然后它阻塞主线程,直到它知道的线程退出。这会延迟解释器的关闭。
- 工作线程对
time.sleep 的调用返回。它调用在其作业的Future 中注册的回调函数,该函数调度另一个作业。
- 与此列表的第 4 步一样,执行程序将作业排队,并启动另一个线程,因为它还没有所需的编号。
- 新线程尝试从内部队列中获取作业,但从步骤 7 中获取了
None 值,这表明它可能已完成。它看到 _shutdown 全局已设置,因此退出。但在此之前,它会在队列中添加另一个 None。
- 第一个工作线程完成其回调。它寻找一个新作业,并找到它在第 8 步中自己排队的那个。它开始运行该作业,就像在第 5 步中一样,这需要一段时间。
- 但没有其他任何事情发生,因为此时第一个工作线程是唯一的活动线程(主线程被阻塞,等待第一个工作线程死亡,而另一个工作线程自行关闭)。
- 我们现在重复步骤 8-12 几次。第一个工作线程将第三个到第 8 个作业排队,并且执行器每次都产生一个相应的线程,因为它没有完整的集合。但是,每个线程都会立即死亡,因为它会从作业队列中获得
None 而不是要完成的实际作业。第一个工作线程最终完成所有实际工作。
- 最后,在第 8 份工作之后,工作方式有所不同。这一次,当回调调度另一个作业时,不会产生额外的线程,因为执行器知道它已经启动了请求的 8 个线程(它不知道 7 个已经关闭)。
- 所以这一次,位于内部作业队列头部的
None 被第一个工作人员(而不是实际作业)拾取。这意味着它会关闭,而不是做更多的工作。
- 当第一个worker关闭时,主线程(一直在等待它退出)终于可以解除阻塞,
_python_exit函数完成。这让解释器完全关闭。我们完成了!
这解释了我们看到的输出!我们得到 8 个输出,全部来自同一个工作线程(第一个产生的)。
我认为在该代码中可能存在竞争条件。如果第 11 步发生在第 10 步之前,事情可能会中断。如果第一个工作人员从队列中获得了None 而另一个新生成的工作人员得到了真正的工作,则将交换角色(第一个工作人员会死,另一个工作人员会做剩下的工作,除非更多的种族这些步骤的后续版本中的条件)。但是,一旦第一个工作人员死亡,主线程就会被解除阻塞。由于它不知道其他线程(因为当它列出要等待的线程时它们并不存在),它会提前关闭解释器。
我不确定这场比赛是否会发生。我猜这不太可能,因为新线程开始和它从队列中获取作业之间的代码路径长度比现有线程完成回调的路径短得多(排队之后的部分)新工作),然后在队列中寻找另一个工作。
我怀疑ThreadPoolExecutor 让我们在将代码作为脚本运行时干净地退出是一个错误。除了执行器自己的self._shutdown 属性之外,排队新作业的逻辑可能还应该检查全局_shutdown 标志。如果是这样,在主线程完成后尝试排队另一个作业会引发异常。
您可以通过在 with 语句中创建 ThreadPoolExecutor 来复制我认为更明智的行为:
# create the pool below the definition of recursive()
with ThreadPoolExecutor(max_workers=8) as pool:
Task(recursive(0)).step()
这将在主线程从step() 调用返回后不久崩溃。它看起来像这样:
exception calling callback for <Future at 0x22313bd2a20 state=finished returned NoneType>
Traceback (most recent call last):
File "S:\python36\lib\concurrent\futures\_base.py", line 324, in _invoke_callbacks
callback(self)
File ".\task_coroutines.py", line 21, in _wakeup
self.step(result)
File ".\task_coroutines.py", line 14, in step
fut = self._gen.send(value)
File ".\task_coroutines.py", line 30, in recursive
Task(recursive(n+1)).step()
File ".\task_coroutines.py", line 14, in step
fut = self._gen.send(value)
File ".\task_coroutines.py", line 28, in recursive
yield pool.submit(time.sleep, 1)
File "S:\python36\lib\concurrent\futures\thread.py", line 117, in submit
raise RuntimeError('cannot schedule new futures after shutdown')
RuntimeError: cannot schedule new futures after shutdown