【问题标题】:How can you wait for completion of a callback submitted from another thread?你怎么能等待从另一个线程提交的回调完成?
【发布时间】:2019-04-06 00:24:31
【问题描述】:

我有两个共享某些状态的 Python 线程,AB。在某一时刻,A 提交了一个回调,由B 在其循环中运行,类似于:

# This line is executed by A
loop.call_soon_threadsafe(callback)

在此之后我想继续做其他事情,但我想确保callback 在这样做之前已经由B 运行。有没有办法(除了标准线程同步原语)让A等待回调完成?我知道call_soon_threadsafe返回一个可以取消任务的asyncio.Handle对象,但是我不确定这是否可以用于等待(我对asyncio还是不太了解)。

在这种情况下,此回调调用loop.close() 并取消剩余的任务,然后在B 中,在loop.run_forever() 之后有一个loop.close()。因此,对于这个用例,特别是一种线程安全机制,它允许我从A 知道循环何时有效关闭也对我有用 - 再次,不涉及互斥体/条件变量/等。

我知道asyncio 并不是线程安全的,只有极少数例外,但我想知道是否提供了一种方便的方法来实现这一点。


这是我的意思的一个非常小的 sn-p,以防万一。

import asyncio
import threading
import time

def thread_A():
    print('Thread A')
    loop = asyncio.new_event_loop()
    threading.Thread(target=thread_B, args=(loop,)).start()
    time.sleep(1)
    handle = loop.call_soon_threadsafe(callback, loop)
    # How do I wait for the callback to complete before continuing?
    print('Thread A out')

def thread_B(loop):
    print('Thread B')
    asyncio.set_event_loop(loop)
    loop.run_forever()
    loop.close()
    print('Thread B out')

def callback(loop):
    print('Stopping loop')
    loop.stop()

thread_A()

我已经用asyncio.run_coroutine_threadsafe 尝试过这种变体,但它不起作用,而是线程A 永远挂起。不知道我做错了什么还是因为我正在停止循环。

import asyncio
import threading
import time

def thread_A():
    global future
    print('Thread A')
    loop = asyncio.new_event_loop()
    threading.Thread(target=thread_B, args=(loop,)).start()
    time.sleep(1)
    future = asyncio.run_coroutine_threadsafe(callback(loop), loop)
    future.result()  # Hangs here
    print('Thread A out')

def thread_B(loop):
    print('Thread B')
    asyncio.set_event_loop(loop)
    loop.run_forever()
    loop.close()
    print('Thread B out')

async def callback(loop):
    print('Stopping loop')
    loop.stop()

thread_A()

【问题讨论】:

标签: python python-asyncio


【解决方案1】:

回调被设置并且(大部分)忘记了。它们不打算用于您需要从中获取结果的东西。这就是为什么生成的句柄只允许您取消回调(不再需要此回调),仅此而已。

如果您需要在另一个线程中等待来自异步管理的协程的结果,请使用协程并将其安排为带有asyncio.run_coroutine_threadsafe() 的任务;这会给你一个Future() instance,然后你可以等待它完成。

但是,使用run_coroutine_threadsafe() 停止循环确实需要循环处理多于它实际能够运行的一轮回调;否则,run_coroutine_threadsafe() 返回的 Future 将不会被通知它计划的任务的状态更改。您可以通过在关闭循环之前在线程 B 中运行 asyncio.sleep(0)loop.run_until_complete() 来解决此问题:

def thread_A():
    # ... 
    # when done, schedule the asyncio loop to exit
    future = asyncio.run_coroutine_threadsafe(shutdown_loop(loop), loop)
    future.result()  # wait for the shutdown to complete
    print("Thread A out")

def thread_B(loop):
    print("Thread B")
    asyncio.set_event_loop(loop)
    loop.run_forever()
    # run one last noop task in the loop to clear remaining callbacks
    loop.run_until_complete(asyncio.sleep(0))
    loop.close()
    print("Thread B out")

async def shutdown_loop(loop):
    print("Stopping loop")
    loop.stop()

当然,这有点 hacky,并且取决于回调管理和跨线程任务调度的内部结构不会改变。正如默认的 asyncio 实现所代表的那样,运行单个 noop 任务对于创建更多正在处理的回调的几轮回调来说已经足够了,但替代循环实现可能会以不同的方式处理这个问题。

所以对于关闭循环,使用基于线程的协调可能会更好:

def thread_A():
    # ...
    callback_event = threading.Event()
    loop.call_soon_threadsafe(callback, loop, callback_event)
    callback_event.wait()  # wait for the shutdown to complete
    print("Thread A out")

def thread_B(loop):
    print("Thread B")
    asyncio.set_event_loop(loop)
    loop.run_forever()
    loop.close()
    print("Thread B out")

def callback(loop, callback_event):
    print("Stopping loop")
    loop.stop()
    callback_event.set()

【讨论】:

  • 我刚刚在阅读有关run_coroutine_threadsafe() 的信息。所以这会给我一个可以安全等待的未来,即使是从另一个线程,即使该任务实际上停止了循环?
  • 我尝试过使用while fut.running(),但这似乎使线程A 立即完成。即使我在callback 中添加time.sleep(1)(仅用于测试!我知道你应该使用asyncio.sleep),future.running() 总是错误的......
  • @jdehesa:我看到的问题是,停止循环也几乎可以保证将信号从管理 shutdown_loop() 协程的任务传递回线程 B 持有的 Future() 的回调是也取消了。连接在这里被切断。关闭循环需要更多……处理。
  • @jdehesa:对此感到抱歉,我已经弄清楚了哪些回调被清除了,只需要一次循环即可。调度 asyncio.sleep(0),一个虚拟 noop 协程,与 asyncio.run_until_complete(),足以清除这些回调。在关闭循环之前在线程 B 上运行它。
  • 是的,这行得通!就像你说的那样,为了清晰和可靠,我最终可能会使用基于线程的协调,但我很高兴找到一种没有它的方法。
【解决方案2】:

是否有任何方法(除了标准线程同步原语)让 A 等待回调完成?

通常你会使用run_coroutine_threadsafe,正如 Martijn 最初建议的那样。但是您对loop.stop() 的使用使回调有些具体。鉴于此,您可能最好使用标准线程同步原语,在这种情况下,这些原语非常简单,并且可以与回调实现和其余代码完全分离。例如:

def submit_and_wait(loop, fn, *args):
    "Submit fn(*args) to loop, and wait until the callback executes."
    done = threading.Event()
    def wrap_fn():
        try:
            fn(*args)
        finally:
            done.set()
    loop.call_soon_threadsafe(wrap_fn)
    done.wait()

不要使用loop.call_soon_threadsafe(callback),而是使用submit_and_wait(loop, callback)。线程同步在那里,但完全隐藏在submit_and_wait中。

【讨论】:

  • 是的,这很好用而且非常简单。我接受了另一个答案,因为它包括一种在没有显式线程同步的情况下解决它的方法,正如问题中所建议的那样,但最终可能也会使用类似的东西。
猜你喜欢
  • 1970-01-01
  • 2021-06-21
  • 2012-03-28
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-06-07
相关资源
最近更新 更多