You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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 Condition variable with an internal lock to ensure all queue operations are thread-safe.

内容的提问来源于stack exchange,提问作者wmIbb

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.13 07:42:04