【发布时间】:2023-02-11 08:15:42
【问题描述】:
我在微服务中有一个类,如下所示:
import asyncio
import threading
class A:
def __init__(self):
self.state = []
self._flush_thread = self._start_flush()
self.tasks = set()
def _start_flush(self):
threading.Thread(target=self._submit_flush).start()
def _submit_flush(self):
self._thread_loop = asyncio.new_event_loop()
self._thread_loop.run_until_complete(self.flush_state()) #
async def regular_func(self):
# This function is called on an event loop that is managed by asyncio.run()
# process self.state, fire and forget next func
task = asyncio.create_task(B.process_inputs(self.state)) # Should call process_inputs in the main thread event loop
self.tasks.add(task)
task.add_done_callback(self.tasks.discard)
pass
async def flush_state(self):
# flush out self.state at regular intervals, to next func
while True:
# flush state
asyncio.run_coroutine_threadsafe(B.process_inputs(self.state), self._thread_loop) # Calls process_inputs in the new thread event loop
await asyncio.sleep(10)
pass
class B:
@staticmethod
async def process_inputs(self, inputs):
# process
在这两个线程上,我有两个独立的事件循环,以避免主事件循环中的任何其他异步函数阻止其他异步函数运行。
我看到 asyncio.run_coroutine_threadsafe 是 thread safe when submitting to a given event loop. 在不同事件循环之间调用的 asyncio.run_coroutine_threadsafe(B.process_inputs()) 仍然是线程安全的吗?
编辑:
process_inputs 将状态上传到对象存储并使用我们传入的状态调用外部 API。
【问题讨论】:
-
在不知道“process_inputs”实际执行和返回的内容的情况下,没有答案。调用
asyncio.run_coroutine_threadsafe(B.process_inputs())在调用线程中执行“process_inputs”,并期望它返回一个协程供另一个循环执行。 -
process_inputs将状态上传到对象存储并使用我们传入的状态调用外部 API。这有帮助吗? -
如果它不返回协同程序,则将其包装在“asyncio.run_coroutine_threadsafe”的调用中是没有意义的。
-
我怀疑你误读了文档。当它说“将协程提交给‘THE’给定的事件循环”时,它指的是作为函数的第二个参数传递的特定事件循环。在你的代码中你根本没有第二个参数,所以这是一个错误。正如@Michael Butscher 指出的那样,第一个参数不是协程,所以这是另一个错误。仅供参考,每个线程最多可以有一个事件循环,因此要求另一个循环运行协程总是意味着由另一个线程执行协程。
-
任何时候你有一个以上的线程都会有线程安全问题,因为它们可以随时相互抢占。您可能需要使用 Lock 或 Condition 对象来保护某些数据结构。 Asyncio 不会导致任何新的线程安全问题,至少我所知道的不会。它也没有解决任何现有的线程问题。在同一个程序中混合使用 asyncio 和线程是没有问题的,只要你注意什么函数可以在什么上下文中运行。
标签: python multithreading thread-safety python-asyncio event-loop