如何优化asyncio锁机制?解决同名请求正常阻塞但异名请求意外被阻塞的问题
Let's break down why your current implementation is blocking all requests regardless of their names, and fix it step by step.
What's Going Wrong?
There are two critical issues in your code:
Lock Dictionary is Local to Each Context Manager Call
Yourlocksdictionary is defined inside theget_lockasync context manager. That means every time you callasync with get_lock(req_name), a brand new emptylocksdict is created. So even identical request names won't actually share the same lock (your observation that same-name requests block might be a side effect of the second issue below).Request Handler is Processing Requests Serially
In_req_handler, you're creating a single task per request and immediately awaiting it withasyncio.gather. This forces the handler to wait for the current request to finish completely before even accepting the next one—so all requests, regardless of name, are blocked behind the previous one.
Optimized Implementation
Here's the fixed code that addresses both issues:
import asyncio from contextlib import asynccontextmanager import json # Assuming logger is defined elsewhere, e.g., logging.getLogger(__name__) logger = ... # Global dictionary to hold locks per request name (shared across all requests) locks = {} @asynccontextmanager async def get_lock(req_name_): logger.info(f"Managing lock for request {req_name_}") # Get or create the lock for this request name lock = locks.setdefault(req_name_, asyncio.Lock()) async with lock: yield # Clean up the lock if no waiters are left (optional but good practice) if len(lock._waiters) == 0: del locks[req_name_] logger.info(f"Removed unused lock for {req_name_}") logger.info(f"Lock released for {req_name_}, active locks count: {len(locks)}") async def handle_lock_request(req_json_): req_name = req_json_.get('req_name') logger.info(f"Attempting to acquire lock for {req_name}") async with get_lock(req_name): logger.info(f"Lock acquired by {req_name}") await _handle_request(req_json_) logger.info(f"Request {req_name} finished processing") async def _handle_request(req_json_): # Replace with your actual request processing logic req_name = req_json_.get('req_name') logger.info(f"Processing request {req_name}") # Simulate work (e.g., API call, database operation) await asyncio.sleep(2) logger.info(f"Request {req_name} finished with action patch stack") async def _req_handler(): # Maintain a list to track active tasks (to avoid memory leaks) active_tasks = [] loop = asyncio.get_running_loop() while True: logger.debug("Waiting for incoming request") req = await receiver.recv_string() # Assuming receiver is a valid async receiver logger.debug(f"Received request: {req}") req_json = json.loads(req) # Create a task for the request and add it to active tasks task = loop.create_task(handle_lock_request(req_json)) active_tasks.append(task) # Clean up completed tasks to prevent memory bloat active_tasks = [t for t in active_tasks if not t.done()] # To run the handler (example) # asyncio.run(_req_handler())
Key Improvements:
- Shared Lock Dictionary: The
locksdict is now global, so all requests share the same set of locks. Thesetdefaultmethod ensures we only create a lock for a request name if it doesn't already exist. - Non-Serial Request Handling: The
_req_handlernow runs in an infinite loop, accepting new requests immediately without waiting for previous ones to finish. Tasks are tracked and cleaned up periodically to avoid memory leaks. - Clean Lock Management: We still remove locks when they have no waiters, keeping the
locksdict lean.
How It Works Now:
- Same-name requests: Will share the same lock, so if one is processing, others will wait until the lock is released.
- Different-name requests: Will use separate locks, so they can run concurrently without blocking each other.
内容的提问来源于stack exchange,提问作者yasin

