【发布时间】:2021-12-30 08:23:07
【问题描述】:
下面是我的代码,我在其中实现了异步锁定机制,如果具有相同名称的请求已经在执行并且尚未完成,那么该机制应该阻止方法中的请求,这工作正常,但问题是如果请求带有不同的名称也被阻止,这并不理想,理想情况下,如果请求带有不同的请求名称,则应该开始执行而无需等待
import asyncio
from contextlib import asynccontextmanager
@asynccontextmanager
async def get_lock(req_name_):
locks = {}
logger.info(f"Creating lock for stack {req_name_} if not created")
if not locks.get(req_name_):
logger.info("creating key for a lock")
locks[req_name_] = asyncio.Lock()
async with locks[req_name_]:
yield
if len(locks[req_name_]._waiters) == 0:
del locks[req_name_]
logger.info(f"lock released")
logger.info(len(locks))
async def handle_lock_request(req_json_):
logger.info(f"ocupying lock")
req_name = req_json_.get('req_name')
async with get_lock(req_name):
logger.info(f"lock acquired by stack {req_name}")
await _handle_request(req_json_)
async def _req_handler():
tasks = []
loop = asyncio.get_running_loop()
logger.debug("Await receiver.recv_string")
req = await receiver.recv_string()
logger.debug(f"Request received {req}")
req_json = json.loads(req)
logger.debug("Await create_task")
tasks.append(loop.create_task(handle_lock_request(req_json)))
await asyncio.gather(*[task for task in tasks if not task.done()])
def _handle_request(req_json_):
# ...
# ...
logger.info(f"Request finished with req name {req_name} for action patch stack")
【问题讨论】:
-
你的
locks字典需要在get_lock()的范围之外。现在,您正在为每个请求创建一个新锁,无论如何。所以我不希望任何请求等待锁定。 -
我尝试了你的建议,通过将 locks dict 保持在 get_lock() 之外效果很好,但我不得不通过删除这一行来对我的代码进行轻微的更改 await asyncio.gather(*[task for task在任务中,如果不是 task.done()]) 因为这一行,即使在将 locks dict 保持在 get_lock() 之外,它也会一个接一个地执行 `
-
你能告诉我这是如何工作的,以防如果获取了锁并且由于某些错误而没有释放它,或者如果请求花费了太多时间来完成,我该如何超时已经运行的请求?你能解释一下超时场景吗
-
您删除的对
gater()的调用不是您的请求按顺序运行的罪魁祸首。当您将其粘贴到您的问题中时,您的代码将无法运行,因为_request_handler()不是异步函数。我现在不知道你的实际代码是什么,但我假设你的request_handler()是一个同步函数,阻塞了你的事件循环。它需要异步或在线程中运行。 -
关于超时请求:请不要在 cmets 中扩大您的问题范围。改为打开一个新问题。
标签: python multithreading python-asyncio locks