Skip to content

Python 多线程

并发(Concurrency)允许多个任务在重叠的时间段内运行或管理执行。Python 提供了几种实现并发的方式,其中线程(Threading)是一种常见的方法。

使用线程的优势:

  • 共享内存:同一进程内的线程共享同一内存空间,与独立进程(多进程 multiprocessing)相比,这使得线程间的通信和数据共享相对直接。
  • 资源效率:线程(通常称为轻量级进程)通常比完整进程消耗更少的系统资源(内存、初始化时间)。
  • 响应性:有助于保持应用程序的响应性(例如,在执行后台任务时处理 GUI 事件),或并发管理多个 I/O 密集型(I/O-bound)操作(如网络请求或文件操作)。

线程(Thread)是可以由调度程序独立管理的最小程序指令序列。一个进程可以包含多个线程,所有线程都在同一进程上下文中执行。

  • 每个线程都有自己的执行上下文(包括指令指针和寄存器)。
  • 线程可以被暂停(进入睡眠状态)或让出控制权,允许其他线程运行。
  • 操作系统调度程序管理线程之间的切换(抢占式多任务处理)。

至关重要的是,在标准的 CPython 解释器中,全局解释器锁 (GIL) 确保在单个进程内,即使在多核处理器上,也只有一个线程同时执行 Python 字节码。这简化了内存管理,但限制了 CPU 密集型(CPU-bound)任务的真正并行性。

GIL 的影响:

  • 线程对于 **I/O 密集型(I/O-bound)**任务是有效的(线程在此类任务中花费时间等待外部操作,如网络或磁盘 I/O)。当一个线程等待时,GIL 可以被释放,允许另一个线程运行。
  • 在 CPython 中,线程对于 **CPU 密集型(CPU-bound)**任务(执行密集计算的任务)不能提供加速,因为任何时候只有一个线程运行 Python 代码。
  • 对于 CPU 密集型任务的并行性,请使用 multiprocessing 模块,它创建独立的进程,每个进程都有自己的解释器和内存空间,从而绕开了 GIL。
  • asyncio 库提供了一种使用事件循环的替代并发模型,适用于高级结构的ated 网络代码和其他 I/O 密集型任务,在管理大量并发连接时通常比线程更有效。

threading 模块是 Python 中用于处理线程的现代、高级接口。它构建在较低级(现已废弃)的 _thread 模块(在 Python 2 中曾称为 thread)之上。新代码始终优先使用 threading 而非 _thread。

  • threading.Thread: 表示执行线程的类。
  • threading.Lock, threading.RLock: 同步原语(synchronization primitives),用于在多个线程访问共享资源时防止竞态条件(race conditions)。
  • threading.Event: 一种简单的线程间通信机制(一个线程发出事件信号,其他线程等待该事件)。
  • threading.Condition, threading.Semaphore, threading.Barrier: 更高级的同步工具。
  • threading.current_thread(): 返回当前的 Thread 对象。
  • threading.active_count(): 返回活动 Thread 对象的数量。

使用 threading.Thread 创建线程主要有两种方式:

1. 将可调用对象(函数)传递给构造函数:

Section titled “1. 将可调用对象(函数)传递给构造函数:”

对于简单的任务,这通常是更直接的方法。

#!/usr/bin/env python3
import threading
import time
def print_time(thread_name, delay):
"""Prints the current time every 'delay' seconds, 5 times."""
count = 0
while count < 5:
time.sleep(delay)
count += 1
current_time = time.strftime("%H:%M:%S", time.localtime())
print(f"{thread_name}: {current_time} (Count: {count})")
print(f"{thread_name} exiting.")
# Create threads by passing the target function and arguments
try:
thread1 = threading.Thread(target=print_time, args=("Thread-1", 1))
thread2 = threading.Thread(target=print_time, args=("Thread-2", 2))
# Start the new threads
thread1.start()
thread2.start()
# Wait for both threads to complete before the main thread exits
thread1.join()
thread2.join()
except Exception as e:
print(f"Error: unable to start thread - {e}")
print("Main thread exiting.")

重写 __init__ 方法(可选,调用父类 __init__)和 run() 方法(必需,包含线程的执行逻辑)。

#!/usr/bin/env python3
import threading
import time
class MyThread(threading.Thread):
def __init__(self, thread_id, name, delay):
# Call the Thread class's initializer
super().__init__()
# Or: threading.Thread.__init__(self)
self.thread_id = thread_id
self.name = name
self.delay = delay
def run(self):
"""This method is executed when thread.start() is called."""
print(f"Starting {self.name}")
# Call the function containing the thread's main logic
print_task(self.name, self.delay)
print(f"Exiting {self.name}")
def print_task(thread_name, delay):
"""A simple task for the thread to perform."""
count = 0
while count < 3:
time.sleep(delay)
count += 1
current_time = time.strftime("%H:%M:%S", time.localtime())
print(f"{thread_name}: {current_time} (Count: {count})")
# Create new threads by instantiating the subclass
thread1 = MyThread(1, "Thread-A", 1)
thread2 = MyThread(2, "Thread-B", 1.5)
# Start new Threads (this calls the run() method)
thread1.start()
thread2.start()
# Wait for threads to complete
thread1.join()
thread2.join()
print("Main thread exiting.")
  • start(): 启动线程活动。它在单独的控制线程中调用 run() 方法(或目标函数)。
  • run(): 表示线程活动的方法。在子类中重写此方法。
  • join([timeout]): 等待线程终止。阻塞调用线程。如果给出 timeout,则最多等待指定秒数。
  • is_alive(): 如果线程仍在执行则返回 True,否则返回 False。
  • name / getName() / setName(): 获取或设置线程的名称。
  • ident: 线程的唯一标识符(如果未启动则为 None)。

当多个线程访问共享资源(如变量或数据结构)时,可能会发生竞态条件(race conditions),导致不可预测的结果。同步原语(synchronization primitives)如锁(locks)用于控制访问。

一个 threading.Lock 提供互斥(mutual exclusion)。一次只能有一个线程“获取”(acquire)锁。其他尝试获取锁的线程将阻塞,直到锁被“释放”(released)。

with 语句通过在进入块时自动获取锁并在退出时释放锁来简化锁的使用,即使发生错误也是如此。

#!/usr/bin/env python3
import threading
import time
shared_counter = 0
counter_lock = threading.Lock() # Create a lock object
class WorkerThread(threading.Thread):
def __init__(self, name):
super().__init__(name=name)
def run(self):
global shared_counter
print(f"{self.name} starting.")
for _ in range(100000):
# Acquire the lock before accessing the shared resource
with counter_lock:
# --- Critical Section Start ---
current_value = shared_counter
# Simulate some processing time
time.sleep(0.000001)
shared_counter = current_value + 1
# --- Critical Section End ---
# Lock is automatically released here
print(f"{self.name} finished.")
# Create threads
thread1 = WorkerThread("Worker-1")
thread2 = WorkerThread("Worker-2")
# Start threads
thread1.start()
thread2.start()
# Wait for threads to complete
thread1.join()
thread2.join()
print(f"Main thread exiting. Final counter value: {shared_counter}")
# Without the lock, the final value would likely be less than 200000 due to race conditions.

你可以手动调用 lock.acquire() 和 lock.release(),但你必须确保 release() 被调用,通常使用 try...finally 块。

# Inside the run method (alternative to 'with')
counter_lock.acquire()
try:
# Access shared resource
current_value = shared_counter
shared_counter = current_value + 1
finally:
counter_lock.release() # Ensure lock is always released

使用队列进行线程通信(queue 模块)

Section titled “使用队列进行线程通信(queue 模块)”

queue 模块(注意在 Python 3 中是小写 ‘q’)提供了线程安全的队列类,非常适合在生产者线程和消费者线程之间传递数据。

  • queue.Queue(maxsize=0): 创建一个先进先出(FIFO)队列。maxsize=0 表示无限大小。
  • put(item, block=True, timeout=None): 将 item 添加到队列。如果 block 为 True 且队列已满(如果 maxsize > 0),则等待直到有空间可用(或 timeout 超时)。
  • get(block=True, timeout=None): 从队列中移除并返回一个项。如果 block 为 True 且队列为空,则等待直到有项可用(或 timeout 超时)。
  • qsize(): 返回队列中项的近似数量。
  • empty(): 如果队列为空则返回 True,否则返回 False。
  • full(): 如果队列已满则返回 True,否则返回 False。
  • task_done(): 在处理完通过 get() 获取的项后调用此方法。与 join() 一起使用。
  • join(): 阻塞直到队列中的所有项都被获取并处理完毕(即,每个项都调用了 task_done())。

其他队列类型:LifoQueue(后进先出)、PriorityQueue(优先队列)。

#!/usr/bin/env python3
import queue
import threading
import time
# A thread-safe queue
work_queue = queue.Queue(10)
exit_flag = False
class ProducerThread(threading.Thread):
def __init__(self, name, q):
super().__init__(name=name)
self.q = q
def run(self):
print(f"{self.name} starting.")
for i in range(1, 6):
if exit_flag:
break
item = f"Data_{i}"
self.q.put(item) # Add item to the queue
print(f"{self.name} produced {item}")
time.sleep(0.5)
print(f"{self.name} finished producing.")
class ConsumerThread(threading.Thread):
def __init__(self, name, q):
super().__init__(name=name)
self.q = q
def run(self):
print(f"{self.name} starting.")
while not exit_flag or not self.q.empty():
try:
# Get item from queue, wait up to 1 sec if empty
item = self.q.get(block=True, timeout=1)
print(f"{self.name} consuming {item}")
# Simulate processing time
time.sleep(1)
self.q.task_done() # Signal that processing is complete
except queue.Empty:
if exit_flag:
break # Exit if flag is set and queue is empty
else:
continue # Queue is temporarily empty, try again
except Exception as e:
print(f"{self.name} encountered error: {e}
")
self.q.task_done() # Still mark task done if error occurs
print(f"{self.name} finished consuming.")
# Create threads
producer = ProducerThread("Producer", work_queue)
consumer1 = ConsumerThread("Consumer-1", work_queue)
consumer2 = ConsumerThread("Consumer-2", work_queue)
# Start threads
producer.start()
consumer1.start()
consumer2.start()
# Wait for producer to finish
producer.join()
# Wait for the queue to be empty (all items processed)
print("Waiting for consumers to finish...")
work_queue.join() # Blocks until task_done() called for all items
# Signal consumers to exit
print("Setting exit flag.")
exit_flag = True
# Wait for consumers to exit cleanly
consumer1.join()
consumer2.join()
print("Main thread exiting.")