Python多线程编程

11255 字
29 分钟

Python多线程编程

发布于

多进程解决了”算”的问题,这一章解决”等”的问题。

上一章用多进程跑满多核,但每个进程独立内存、通信靠序列化,
对 I/O 密集型任务来说过于沉重。网络请求、数据库查询这类场景,
CPU 大部分时间在空转,瓶颈不在算力,在等待。真正需要的不是更多核心,
而是”谁在等,就让别人先干”。

多线程正是为此而生:共享内存、创建轻量、切换高效。
不过线程并不是终点——当成百上千个请求同时到来,线程本身的开销
也会成为新的瓶颈。但在讨论那个更轻量的方案之前,我们先打好基础。

前置知识

  • Python 函数与类的基础用法

核心内容

线程基础概念

进程 VS 线程

  • 进程:操作系统资源分配的最小单位。每个进程拥有独立的内存地址空间、 文件描述符表和系统资源,进程间通信需要通过管道、共享内存等机制, 开销较大。

  • 线程:进程内 CPU 调度的最小单位。同一进程内的线程共享内存空间和 文件描述符等资源,创建和切换的开销远小于进程,又称轻量级进程。

  • 关系:一个进程至少包含一个主线程,可动态创建多个子线程。 同一进程内的所有线程共享全局变量和堆内存,但各自拥有独立的栈和 程序计数器。

进程线程
资源分配独立内存、文件描述符表共享所属进程资源
创建开销重(MB 级)轻(KB 级)
并行能力真并行,绕过 GIL受 GIL 限制
适合场景CPU 密集型I/O 密集型

GIL(Global Interpreter Lock)

GIL是CPython 中的一把全局锁,保证同一时刻只有一个线程执行 Python 字节码。 这意味着多线程无法真正利用多核 CPU 进行并行计算。所以可以:

时间线(两个线程共享一个 GIL)
──────────────────────────────────────────
线程A: [计算] [发起网络请求 → 让出GIL] [计算] [完成]
线程B:         [  获得GIL → 趁机计算  ]
──────────────────────────────────────────

创建线程的两种方式

方式1:函数式

import threading
import time

def task(name, delay):
    """模拟耗时任务"""
    print(f"线程 {name} 启动")
    time.sleep(delay)
    print(f"线程 {name} 执行完毕")

if __name__ == "__main__":
    start = time.perf_counter()

    # 创建两个线程
    t1 = threading.Thread(target=task, args=("A", 2))
    t2 = threading.Thread(target=task, kwargs={"name": "B", "delay": 1})

    # 启动线程
    t1.start()
    t2.start()

    # 等待两个线程都结束
    t1.join()
    t2.join()

    elapsed = time.perf_counter() - start
    print(f"所有子线程执行完成,耗时 {elapsed:.1f}s")
    # 串行需要 3s,并行只需 2s

方式2:类继承 Thread

import threading
import time

class MyThread(threading.Thread):
    def __init__(self, name, delay):
        super().__init__(name=name)
        self.delay = delay

    def run(self):
        print(f"线程 {self.name} 启动")
        time.sleep(self.delay)
        print(f"线程 {self.name} 结束")

if __name__ == "__main__":
    t1 = MyThread("线程1", 1)
    t2 = MyThread("线程2", 2)
    t1.start()
    t2.start()
    t1.join()
    t2.join()
    print("程序结束")                # 逐个等待,全部跑完再继续

join() 不能省

  1. 主线程/主进程的代码执行完毕后,Python 解释器不会立即退出。
  2. 解释器会检查是否还有非守护的子线程和子进程在运行。
  3. 有 → 解释器继续等待;
    没有 → 解释器退出,所有守护线程和守护子进程被强制销毁。

简记:所有非守护的子线程和子进程都结束,程序才真正退出。
join() 阻塞调用线程(通常是主线程),直到被 join 的线程终止,或超时。
❌ 没有 join —— 可能在子线程结束前就用了结果;
✅ 有 join —— 确保拿到结果。

补充:守护线程

import threading
import time

def daemon_task():
    while True:
        print('守护线程正在运行中...')
        time.sleep(1)

def normal_task():
    time.sleep(3)
    print('普通子线程结束')

if __name__ == '__main__':
    # daemon=True,说明 t1 是守护线程,主线程结束,t1 直接被回收
    t1 = threading.Thread(target=daemon_task, daemon=True)
    # t2 不是守护线程,主线程要等待 t2 结束之后才能结束
    t2 = threading.Thread(target=normal_task)

    t1.start()
    t2.start()
    t2.join()
    print('主线程执行完毕,守护线程直接销毁...')

运行逻辑:3 秒后普通线程结束,主线程结束,无限循环的守护线程直接终止。

daemon 必须在 start() 前设置,启动后修改无效

守护线程不能执行文件、数据库落盘操作,主线程退出会强制杀死,数据丢失。

补充:线程常用属性与方法

方法 / 属性作用
start()启动线程,调用 run ()
run()线程执行逻辑入口
join(timeout=None)阻塞主线程等待子线程,timeout 设置最大等待秒数
is_alive()返回布尔值,判断线程是否正在运行
getName() / name获取 / 设置线程名称
threading.current_thread()获取当前运行线程对象
threading.enumerate()返回当前所有存活线程列表
threading.active_count()获取当前活跃线程总

线程同步工具

多个线程同时改同一个变量会发生竞态:

import threading
import time

num = 0

def add():
    global num
    for _ in range(100_000):
        current = num
        time.sleep(0)  # 不睡,但让出 GIL,立刻重新排队
        num = current + 1  # 非原子操作,多线程互相覆盖 → 结果偏小

if __name__ == "__main__":
    t1 = threading.Thread(target=add)
    t2 = threading.Thread(target=add)
    t1.start()
    t2.start()
    t1.join()
    t2.join()
    print("预期200_000,实际结果:", num)    # 期望 200_000,实际通常不是

互斥锁 Lock

同一时间仅一个线程获取锁,保证代码块原子操作。

lock = threading.Lock()

def add():
    global num
    for _ in range(100_000):
        with lock:        # 等价于 acquire() + try/finally + release()
            current = num
            time.sleep(0)
            num = current + 1

锁的使用原则
with lock: 包裹的代码(锁的范围)越小越好,但不能小到漏掉需要保护的操作。

锁保护的是数据,不是代码。
共享资源要锁,独立计算不锁。
检查和修改必须在同一把锁内完成。

可重入锁 RLock

函数嵌套调用时需要重复获取锁,普通 Lock 同一线程重复 acquire 会死锁,RLock 支持同一线程多次加锁,需要同等次数释放:

import threading

rlock = threading.RLock()

def transfer(account_from, account_to, amount):
    with rlock:
        # 获取锁后,调用内部方法
        deduct(account_from, amount)
        deposit(account_to, amount)

def deduct(account, amount):
    with rlock:      # 同一线程再次获取,RLock 允许
        account.balance -= amount

def deposit(account, amount):
    with rlock:      # 同一线程再次获取,RLock 允许
        account.balance += amount

死锁规避方案

  1. 统一所有线程获取锁的顺序;
  2. 设置 acquire 超时 lock.acquire(timeout=3);
  3. 减少嵌套锁;
  4. 使用 RLock 替代普通 Lock。

信号量 Semaphore:限制最大并发线程

控制同时运行线程数量:

'''
    Lock           → 独占,一次 1 个      → 保护共享资源
    Semaphore(n)   → 共享,一次 n 个      → 限制并发数量
    RLock          → 可重入,同一线程多次获取 → 函数嵌套调用
'''
import threading
import time

# 同时允许 3 个线程访问
sem = threading.Semaphore(3)

def access_resource(thread_id):
    print(f"线程{thread_id}: 等待进入...")
    with sem:
        print(f"线程{thread_id}: ✅ 进入,正在使用资源")
        time.sleep(2)
        print(f"线程{thread_id}: ❌ 离开")

if __name__ == "__main__":
    threads = []
    for i in range(8):
        t = threading.Thread(target=access_resource, args=(i,))
        threads.append(t)
        t.start()

    for t in threads:
        t.join()
    print("全部完成")

事件 Event:线程间通知机制

核心方法:

  • event.set():设置标志为 True,唤醒等待线程
  • event.wait(timeout):阻塞等待标志为 True
  • event.clear():重置标志为 False
import threading
import time

event = threading.Event()

def wait_task():
    print('子线程等待信号...')
    event.wait()    # 阻塞,直到 event 被 set
    print('收到信号,并且开始执行任务')

def send_signal():
    time.sleep(3)
    print('发送通知信号')
    event.set()    # 通知所有等待的线程

if __name__ == '__main__':
    t1 = threading.Thread(target=wait_task)
    t2 = threading.Thread(target=send_signal)
    t1.start()
    t2.start()
    t1.join()
    t2.join()

条件变量 Condition:复杂生产消费模型

Condition = 锁 + 等待队列。线程在条件不满足时 wait() 等待,条件满足时被 notify_all() 唤醒。
适合「生产者-消费者」这类一写多读的复杂模型。

wait() / notify_all() 机制

wait() 三步曲:释放锁 → 进入等待队列 → 阻塞
notify_all():唤醒等待队列 → 回到锁池重新抢锁

被唤醒 ≠ 立即执行:抢到锁的线程才从 wait() 的下一行继续

锁池 vs 等待队列

锁池 :等待拿锁的线程,参与锁竞争
等待队列 :已 wait() 的线程,不参与竞争
wait() → 进等待队列;notify_all() → 回到锁池抢锁

完整可运行示例:

import threading
import time
import random

BUFFER = []                    # 共享仓库
MAX, TOTAL = 5, 20             # 仓库容量 / 总任务数
produced = consumed = 0
done = False
cond = threading.Condition()   # 锁 + 等待队列

def producer(name):
    global produced, done
    while True:
        with cond:
            if produced >= TOTAL:
                done = True
                cond.notify_all()
                return
            while len(BUFFER) >= MAX and not done:
                cond.wait()             # 仓库满 → 等待
            if done:                    # 第二次检查,防止过度生产
                cond.notify_all()
                return
            produced += 1
            BUFFER.append(f"{name}-P{produced}")
            cond.notify_all()           # 唤醒消费者
        time.sleep(random.uniform(0.1, 0.3))

def consumer(name):
    global consumed, done
    while True:
        with cond:
            while not BUFFER and not done:
                cond.wait()             # 仓库空 → 等待
            if not BUFFER:              # 没货且生产结束 → 退出
                done = True
                cond.notify_all()
                return
            item = BUFFER.pop(0)
            consumed += 1
            if consumed >= TOTAL:
                done = True
                cond.notify_all()
                return
            cond.notify_all()
        time.sleep(random.uniform(0.1, 0.3))

if __name__ == "__main__":
    ts = [threading.Thread(target=producer, args=(f"P{i+1}",)) for i in range(2)] + \
         [threading.Thread(target=consumer, args=(f"C{i+1}",)) for i in range(3)]
    for t in ts:
        t.start()
    for t in ts:
        t.join()
    print(f"完成!总生产 {produced},总消费 {consumed}")

为什么用while 不是 if?
被唤醒 ≠ 条件一定满足:可能有虚假唤醒,或货已被其他消费者抢走。必须 while 重新检查,否则会越界 pop。

惊群效应:notify_all() 唤醒全部线程,但只有抢到锁的那个能干活,其余又回去等待,浪费 CPU。单个唤醒用 notify(),但可能造成线程饥饿。

线程本地存储 threading.local ()

多线程共享全局变量会引发数据竞争。threading.local() 为每个线程 创建独立的变量副本,线程之间互不干扰,无需加锁。

import threading

# 共享变量 —— 会竞争
counter = 0

# 线程私有变量 —— 互不干扰
local = threading.local()

def task(name):
    global counter

    # 每个线程有自己的 local.value
    local.value = 0
    for _ in range(10000):
        local.value += 1
    print(f"{name} 的私有值: {local.value}")

    # 共享变量 —— 有竞争
    for _ in range(10000):
        counter += 1

if __name__ == "__main__":
    threads = [threading.Thread(target=task, args=(f"线程{i}",)) for i in range(3)]
    for t in threads:
        t.start()
    for t in threads:
        t.join()

    print(f"\n私有变量: 各线程互不影响,结果都是 10000")
    print(f"共享变量: counter = {counter}(预期 30000,实际不足)")

补充:线程间通信方式汇总表

机制作用适合场景
Lock互斥,一次一个线程保护共享资源
RLock可重入锁函数嵌套调用同一把锁
Semaphore限制并发数连接池、限流
Event一对多通知启动信号、状态广播
Condition精准唤醒生产消费模型
Queue线程安全的队列生产消费(推荐替代 Condition)
local线程私有变量避免共享、无需加锁

线程池:ThreadPoolExecutor

手写线程需要管理 start、join、线程生命周期,代码冗余且容易出错。
ThreadPoolExecutor 提供了线程池封装,自动管理线程的创建、复用和回收。

类型导入适用场景并发机制
ThreadPoolExecutorconcurrent.futures.ThreadPoolExecutorI/O 密集型(网络、文件)多线程,受 GIL 限制
ProcessPoolExecutorconcurrent.futures.ProcessPoolExecutorCPU 密集型(计算、数据)多进程,绕过 GIL

基本用法

from concurrent.futures import ThreadPoolExecutor
import time

def task(num):
    print(f"任务{num}开始")
    time.sleep(1)
    return f"任务{num}完成"

if __name__ == "__main__":
    with ThreadPoolExecutor(max_workers=4) as pool:   # max_workers:最大并发数
        # 写法一:submit —— 返回 Future,用 .result() 取值
        futures = [pool.submit(task, i) for i in range(6)]
        for f in futures:
            print(f.result())

        # 写法二:map —— 直接返回结果迭代器
        for r in pool.map(task, range(6)):
            print(r)

ThreadPoolExecutor 提供两种提交任务的方式:

写法返回值取结果适合场景
submitFuture 对象.result()需要精细控制(取消、回调、超时)
map结果迭代器直接遍历同一函数批量调用

两种写法执行流程完全一致,只是封装层级不同:submit 手动创建每个 Future 再 .result() 取值;map 内部自动完成创建与收集。

等待机制:结果什么时候到

submit 拿到 Future 后,取结果有三种方式,差别在「顺序 vs 及时性 vs 精确控制」:

方式结果顺序一句话适用场景
遍历 .result()按提交顺序逐个阻塞等任务少、不在乎等待
as_completed()按完成顺序先完成先处理大量 I/O,想尽快消费先到的结果
wait()不取结果只等一个「里程碑」阶段同步

容易踩的坑:for f in futures: f.result() 严格按提交顺序等 —— 如果第 1 个任务最慢,后面全完成了也得干等它。要「先完成先处理」,用 as_completed:

from concurrent.futures import ThreadPoolExecutor, as_completed
import time

def task(num):
    time.sleep(3 - num % 3)          # 让不同任务耗时不同
    return f"任务{num}"

if __name__ == "__main__":
    with ThreadPoolExecutor(max_workers=4) as pool:
        futures = [pool.submit(task, i) for i in range(6)]
        for f in as_completed(futures):   # 谁先完成先输出谁
            print(f.result())

wait() 不直接返回结果,而是等待一组任务达到指定状态后, 把”已完成”和”未完成”分开返回,由你决定下一步怎么处理。

条件场景思路
FIRST_COMPLETED竞速请求、搜索建议谁快用谁
ALL_COMPLETED批量处理、分页聚合全部拿到再继续
FIRST_EXCEPTION批量执行、快速失败有一个错就停
用法参考:
from concurrent.futures import ThreadPoolExecutor, wait, FIRST_COMPLETED, ALL_COMPLETED, FIRST_EXCEPTION
import time
import random

# ============================================================
# FIRST_COMPLETED:谁先完成就先返回
# ============================================================
def fast_search(source):
    time.sleep(random.uniform(0.5, 2))
    return f"来自 {source} 的结果"

print("=== FIRST_COMPLETED ===")
with ThreadPoolExecutor(3) as pool:
    futures = [pool.submit(fast_search, s) for s in ["镜像A", "镜像B", "镜像C"]]
    done, pending = wait(futures, return_when=FIRST_COMPLETED)

    result = list(done)[0].result()
    print(f"最快返回:{result}")
    print(f"还在跑:{len(pending)} 个")

# ============================================================
# ALL_COMPLETED:全部完成才返回
# ============================================================
def process(order):
    time.sleep(random.uniform(0.5, 1.5))
    return f"订单 {order} 已处理"

print("\n=== ALL_COMPLETED ===")
with ThreadPoolExecutor(3) as pool:
    futures = [pool.submit(process, o) for o in ["001", "002", "003"]]
    done, pending = wait(futures, return_when=ALL_COMPLETED)

    print(f"全部完成,共 {len(done)} 个结果:")
    for f in done:
        print(f"  {f.result()}")

# ============================================================
# FIRST_EXCEPTION:有异常就立刻返回
# ============================================================
def risky_task(i):
    time.sleep(0.5)
    if i == 2:
        raise ValueError(f"任务 {i} 失败了!")
    return f"任务 {i} 成功"

print("\n=== FIRST_EXCEPTION ===")
with ThreadPoolExecutor(4) as pool:
    futures = [pool.submit(risky_task, i) for i in range(5)]
    done, pending = wait(futures, return_when=FIRST_EXCEPTION)

    for f in done:
        if f.exception():
            print(f"发现异常:{f.exception()}")
            for p in pending:
                p.cancel()
            break
        else:
            print(f"已完成:{f.result()}")
    print(f"被取消:{len(pending)} 个")

异常捕获 + 任务超时

Future.result (timeout = 秒数) 超时抛 TimeoutError:

from concurrent.futures import ThreadPoolExecutor, TimeoutError
import time

def err_task():
    time.sleep(2)
    raise ValueError("任务内部报错")

with ThreadPoolExecutor(2) as pool:
    f = pool.submit(err_task)
    try:
        res = f.result(timeout=1)
    except TimeoutError:
        print("任务执行超时")
    except Exception as e:
        print("任务异常:", e)

常见问题及补充

什么时候用?

场景结论原因
大量网络请求 / 文件读写多线程CPU 大部分时间在等 I/O,线程等待时不阻塞其他线程
纯计算(循环、图像处理)多进程绕过 GIL,真正利用多核 CPU 并行计算
海量连接的网络服务asyncio单线程协程,万级连接只需极少内存开销
任务少且很快单线程并发本身有开销,任务太轻反而更慢
I/O + 计算混合多线程 + 多进程I/O 用线程,计算提交给进程池

死锁

四个条件:互斥、持有并等待、不可抢占、循环等待。

预防方法:

  1. 统一加锁顺序;
  2. 使用 RLock(解决同一线程嵌套获取);
  3. 减少嵌套锁;
  4. 设置超时(lock.acquire(timeout=5));

实战示例:多线程网络请求对比测试

import time
import urllib.request
from concurrent.futures import ThreadPoolExecutor

URLS = [
    "https://httpbin.org/delay/2",
    "https://httpbin.org/delay/3",
    "https://httpbin.org/delay/1",
]

def fetch(url):
    try:
        with urllib.request.urlopen(url, timeout=5) as resp:
            data = resp.read()
            return url, len(data), None
    except Exception as e:
        return url, 0, str(e)

def show_results(label, results):
    for url, size, err in results:
        status = f"{size} bytes" if not err else f"失败: {err}"
        print(f"  {url} → {status}")

if __name__ == "__main__":
    # 顺序执行
    start = time.perf_counter()
    seq_results = [fetch(u) for u in URLS]
    seq_time = time.perf_counter() - start

    # 多线程并发
    start = time.perf_counter()
    with ThreadPoolExecutor(max_workers=3) as pool:
        con_results = list(pool.map(fetch, URLS))
    con_time = time.perf_counter() - start

    print(f"顺序执行 ({seq_time:.2f}s):")
    show_results("顺序", seq_results)

    print(f"\n并发执行 ({con_time:.2f}s):")
    show_results("并发", con_results)

    print(f"\n加速比: {seq_time / con_time:.2f}x")

结语

多线程擅长 I/O 密集型任务,多进程擅长 CPU 密集型任务, 但两者都有一个共同的限制:当并发量上升到成百上千时, 线程和进程本身的开销也会成为瓶颈。

下一章我们将学习协程——一种更轻量的并发方案, 用单线程实现万级并发。

参考资料

Last updated on