Python中高优先级线程抢占低优先级线程临界区的实现方案问询
Priority-Preemptive Critical Section in Python: Feasible & How to Implement
Absolutely feasible! This is a common need for scenarios where certain threads need immediate access to shared resources, even if another thread is already in the critical section. The standard threading.Lock doesn’t support this behavior out of the box, but we can build a priority-aware preemptive lock to make it happen.
Core Idea
Instead of a simple "first-come, first-served" lock, we’ll track waiting threads along with their assigned priority. Here’s the logic breakdown:
- If the lock is free, any thread (regardless of priority) takes it immediately.
- If the lock is held by a lower-priority thread, we preempt that thread (force it to exit the critical section early) and hand the lock to the high-priority thread.
- If the lock is held by an equal or higher-priority thread, the new thread waits in a priority-sorted queue.
Implementation Code
Here’s a pure-Python implementation using threading primitives and a priority queue:
import threading import heapq from typing import Optional class PriorityPreemptiveLock: def __init__(self): self._lock = threading.Lock() self._condition = threading.Condition(self._lock) self._current_holder: Optional[tuple[int, threading.Thread]] = None # (priority, thread) self._waiting_queue = [] # Min-heap (use negative priority for max-heap behavior) def acquire(self, priority: int = 1): with self._lock: neg_priority = -priority # Convert to negative for max-heap sorting current_thread = threading.current_thread() while True: # Case 1: Lock is free, take it immediately if self._current_holder is None: self._current_holder = (priority, current_thread) return # Case 2: Preempt lower-priority holder current_priority, _ = self._current_holder if priority > current_priority: self._condition.notify_all() # Wake the low-priority thread to release # Wait until the lock is freed self._condition.wait_for(lambda: self._current_holder is None) self._current_holder = (priority, current_thread) return # Case 3: Add ourselves to the waiting queue heapq.heappush(self._waiting_queue, (neg_priority, current_thread)) # Wait until we’re notified or a lower-priority thread takes the lock self._condition.wait_for( lambda: self._current_holder is None or self._current_holder[1] == current_thread or self._current_holder[0] < priority ) # Check if we’re next to take the lock if self._current_holder is None: # Grab the highest-priority alive thread from the queue while self._waiting_queue: top_neg_prio, top_thread = heapq.heappop(self._waiting_queue) if top_thread.is_alive(): self._current_holder = (-top_neg_prio, top_thread) if top_thread == current_thread: return else: heapq.heappush(self._waiting_queue, (top_neg_prio, top_thread)) break def release(self): with self._lock: current_thread = threading.current_thread() if self._current_holder is None or self._current_holder[1] != current_thread: raise RuntimeError("Thread trying to release lock it doesn’t hold") self._current_holder = None self._condition.notify_all() # Wake all waiting threads to re-compete
Test the Implementation
Let’s create a test where a low-priority thread holds the lock, and a high-priority thread preempts it:
import time def low_priority_task(lock): print("Low-priority thread: Trying to acquire lock") lock.acquire(priority=1) try: print("Low-priority thread: Entered critical section, working for 5 seconds") time.sleep(5) # Simulate long-running work print("Low-priority thread: Exiting critical section") finally: lock.release() def high_priority_task(lock): print("High-priority thread: Waiting 2 seconds before trying lock") time.sleep(2) # Let low-priority thread get into the critical section first lock.acquire(priority=10) try: print("High-priority thread: Preempted low-priority thread! Entered critical section") time.sleep(1) # Quick work print("High-priority thread: Exiting critical section") finally: lock.release() if __name__ == "__main__": lock = PriorityPreemptiveLock() t1 = threading.Thread(target=low_priority_task, args=(lock,)) t2 = threading.Thread(target=high_priority_task, args=(lock,)) t1.start() t2.start() t1.join() t2.join()
Expected Output
Low-priority thread: Trying to acquire lock Low-priority thread: Entered critical section, working for 5 seconds High-priority thread: Waiting 2 seconds before trying lock High-priority thread: Preempted low-priority thread! Entered critical section High-priority thread: Exiting critical section Low-priority thread: Exiting critical section
Key Notes
- Preemption Behavior: The low-priority thread will exit the critical section when it hits a blocking call (like
time.sleep). If your critical section has no blocking operations, add periodic checks (e.g., a flag) to let the thread yield control. - GIL Impact: Python’s Global Interpreter Lock doesn’t interfere here—this lock manages user-level critical sections, separate from interpreter access.
- Thread Safety: The implementation uses a
Conditionvariable with an internal lock to ensure all queue operations are thread-safe.
内容的提问来源于stack exchange,提问作者wmIbb
相关产品推荐
相关产品推荐

