【问题标题】:How can I improve asyncio lock mechanism如何改进异步锁机制
【发布时间】: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


【解决方案1】:

更新 从代码中删除 asyncio.gather() 部分,然后调用 create_task

import asyncio
from contextlib import asynccontextmanager

locks = {}
    @asynccontextmanager
    async def get_lock(req_name_):
        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")
          loop.create_task(handle_lock_request(req_json))
          
  
     def _handle_request(req_json_):
       ...
       ...
       logger.info(f"Request finished with req name {req_name} for action patch stack")

【讨论】:

  • 正如目前所写,您的答案尚不清楚。请edit 添加其他详细信息,以帮助其他人了解这如何解决所提出的问题。你可以找到更多关于如何写好答案的信息in the help center
  • 问题和答案对我来说都不清楚。问题需要改进,但首先问题是什么?答案只是很多代码,解决了哪个问题,为什么?代码也没有记录,应该分成更小的例程,每个函数都有一个单独的功能和职责。给出正确的函数名称:handle_something() 对我来说不是很清楚而且太笼统。然后,添加文档,以便您了解它是如何工作的。用几行代码告诉用户何时使用此代码 m.
猜你喜欢
  • 2021-03-14
  • 1970-01-01
  • 2014-11-13
  • 2014-09-12
  • 1970-01-01
  • 1970-01-01
  • 2017-01-05
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多