第一次看到 threading.Thread 时,容易把它理解成“让函数去后台执行”。这个说法只描述了表面现象,没有解释后面的几个问题:为什么调用 start() 后输出顺序变了,为什么有时需要 join(),为什么多个线程修改同一个变量会出错,以及 Queue.get() 为什么可以一直等待。
这些问题都来自同一个基础:Thread 是进程中的另一条执行路径。多条线程共享进程内存,各自推进任务,由操作系统决定何时切换。start()、join()、Lock 和 Queue,分别对应启动、等待、共享数据保护和线程通信。
从串行执行到两条执行路径
没有显式创建线程时,Python 代码看起来是一行接一行执行:
print("A")
print("B")
print("C")A
↓
B
↓
C这段代码并非“没有线程”,而是全部运行在 Python 程序启动时创建的主线程(MainThread)中:
import threading
print(threading.current_thread().name)MainThread进程(Process)可以先理解成程序运行时的资源容器,里面放着代码、内存、打开的文件和 Socket。线程(Thread)是进程中执行代码的路径。一个进程至少有一条主线程,也可以再创建多条工作线程:
Operating System
└── Python Process
├── 代码、对象、文件描述符、Socket
│
├── MainThread
├── Worker Thread A
└── Worker Thread B同一进程中的线程看到的是同一批 Python 对象,并不是每条线程复制一份变量。线程的创建成本通常低于进程,代价是共享状态需要协调。
下面的完整例子把两个等待任务放到两条工作线程中:
import threading
import time
def download(name: str, seconds: int) -> None:
print(f"{name} start: {threading.current_thread().name}")
time.sleep(seconds)
print(f"{name} end")
threads = [
threading.Thread(target=download, args=("A", 2), name="download-A"),
threading.Thread(target=download, args=("B", 2), name="download-B"),
]
for thread in threads:
thread.start()
print("两条工作线程已经启动")
for thread in threads:
thread.join()
print("所有下载任务结束")两次 sleep(2) 会重叠,总耗时接近 2 秒,而不是串行执行时的 4 秒。这里并没有三份 Python 程序,结构仍然是一个进程,只是进程中有三条执行路径:
Python Process
├── MainThread
├── download-A
└── download-B这三条路径不再组成一条固定的源码顺序。主线程启动 A 后可能马上启动 B,也可能在中间让 A 先执行一段;具体先后由调度时机决定。唯一确定的是:每条线程内部仍按自身的代码顺序执行,例如 download-A 一定先打印自己的 start,再打印自己的 end。
time.sleep() 代表等待 I/O 的阶段。线程 A 等待时,线程 B 可以继续推进。真实代码中的 HTTP 请求、数据库访问、文件读写和 Socket 接收也有类似的等待阶段。
Thread、start() 与 join()
Python 程序启动时已经存在一条主线程,可以通过 threading.current_thread() 取得。创建 Thread 对象只是描述要执行的工作,还没有启动新的执行流:
thread = threading.Thread(
target=download,
args=("A", 2),
kwargs={},
)这时的状态可以拆成两层:
Python Process
└── MainThread 正在运行
thread 只是一个 Thread 对象
download() 还没有执行
新的执行路径也还没有启动target 接收可调用对象,因此写的是 target=download。函数名和函数调用表达式在这里含义不同:
| 写法 | 发生的事 |
|---|---|
target=download | 把函数对象交给 Thread,等待线程启动后调用 |
target=download("A", 2) | 先在当前线程立刻调用函数,再把返回值交给 Thread |
如果 download() 没有显式返回值,第二种写法最后近似变成 Thread(target=None)。耗时任务已经在主线程里执行完,新线程反而没有任务可做。
args=("A", 2) 会在线程内部形成一次 download("A", 2) 调用。只有一个位置参数时要写成 args=("A",);逗号决定它是单元素 tuple,而不是带括号的字符串。关键字参数则通过 kwargs 传入:
thread = threading.Thread(
target=download,
kwargs={"name": "A", "seconds": 2},
)这等价于工作线程启动后执行 download(name="A", seconds=2)。
调用 start() 后,Thread 会在另一条执行路径中调用 target。start() 不等待目标函数结束,当前线程会继续向下执行:
MainThread Worker Thread
│ │
├── thread.start() ─────────────>│ download(...)
│ │
├── 继续执行后续代码 │ 继续执行任务
│ │start() 的职责不是直接在当前位置执行 download(),而是安排新线程调用它:
1. MainThread 调用 thread.start()
2. Python 与操作系统启动新的线程
3. Worker Thread 进入 Thread.run()
4. Thread.run() 调用 target(*args, **kwargs)
5. MainThread 同时从 start() 后面继续执行所以这段代码只保证打印过 A 和 B,不保证两者顺序:
def worker():
print("A")
thread = threading.Thread(target=worker)
thread.start()
print("B")可能的输出有两种:
A B
B 或 Astart() 只建立并发关系,不承诺哪条线程先拿到执行机会。连续几次运行都得到 A B,也不能据此推导以后都保持这个顺序。
join() 阻塞的是调用它的线程,直到目标线程结束:
thread.start()
thread.join()
print("B")这里可以确定 A 先于 B。join() 最容易误解的地方是“谁在等谁”:
MainThread 调用 worker_thread.join()
│
├── MainThread 暂停在这里
│
└── worker_thread 继续运行,不受影响
│
└── target 结束
│
└── MainThread 从 join() 后继续它没有把已经开始的并发执行变回串行;工作线程在 start() 后独立运行,只是主线程在需要结果的位置等待。如果主线程在启动每条线程后立刻 join(),并发会被无意间消掉:
# A 完成后才会启动 B,总时间接近 4 秒
for thread in threads:
thread.start()
thread.join()先启动全部线程,再逐个等待,两个任务才有机会重叠:
for thread in threads:
thread.start()
for thread in threads:
thread.join()带超时的 join(timeout=1) 始终返回 None,超时后要通过 thread.is_alive() 判断线程是否仍在运行:
thread.join(timeout=1)
if thread.is_alive():
print("一秒后任务仍未结束")同一个 Thread 对象只能 start() 一次。线程尚未启动时调用 join(),或者线程等待自己结束,都会抛出 RuntimeError。这些行为由 threading 官方文档定义。
并发适合解决什么问题
并发表示多个任务可以在同一段时间内交错推进;并行表示多个任务在同一时刻真正执行。普通 CPython 构建受到 GIL 限制时,多条线程仍能并发,但纯 Python 代码通常不能依靠它们同时占用多个 CPU 核心。
串行执行两个等待任务时,时间线是首尾相接的:
0s 2s 4s
│──── A 等待 ────│──── B 等待 ────│放入两条线程后,两段等待可以重叠:
0s 2s
│──── A 等待 ────│
│──── B 等待 ────│缩短的是“等待被串起来”的时间。线程并没有让每一次网络请求本身从 2 秒变成 1 秒,也没有消除远端服务的延迟。
一个 I/O 任务通常包含少量 CPU 工作和较长等待:
发送 HTTP 请求 1 ms
等待远端响应 500 ms
解析响应内容 2 ms当线程进入网络、文件或 Queue 等待时,调度器可以让其他线程继续。这样的重叠正是 Thread 常见的价值来源。
| 任务 | Thread 的作用 |
|---|---|
| HTTP、数据库、Redis、文件和 Socket I/O | 等待期间让其他线程推进 |
| Queue 消费、后台监听、流式响应 | 将等待和数据处理拆成不同执行路径 |
| 大量纯 Python 计算 | 普通 CPython 构建下通常不能获得多核并行 |
因此,线程常用来减少 I/O 等待造成的串行停顿,而不是默认用来加速所有函数。CPU 密集型任务还要考虑 multiprocessing、ProcessPoolExecutor、能释放 GIL 的原生库,或 free-threaded Python。
传统 CPython 使用全局解释器锁(GIL)。在同一个进程中,通常只有持有 GIL 的线程能执行 Python 字节码;阻塞 I/O 期间会释放 GIL,所以 I/O 密集型任务仍能从多线程并发中受益。
这个结论需要带上版本和构建方式。CPython 从 3.13 开始提供可选的 free-threaded 构建,可以关闭 GIL,让线程在多个 CPU 核心上并行执行 Python 代码,但它不是普通安装的默认模式,部分不兼容的扩展模块还可能在运行时重新启用 GIL。当前解释器是否支持或启用了 free threading,可以按 Python free-threading 官方说明检查。
不论 GIL 是否启用,共享可变状态都需要明确的同步设计。free-threaded 构建改变了可并行执行的能力,没有改变竞态条件本身。
共享内存、竞态条件与 Lock
同一进程中的线程可以访问同一个 list、dict 或业务对象。这让数据共享很直接,也意味着多个线程可能在未协调的情况下修改同一份状态。
最直接的共享例子是一条工作线程修改 list,主线程在等待它结束后读取同一个 list:
import threading
messages: list[str] = []
def worker() -> None:
messages.append("hello")
thread = threading.Thread(target=worker)
thread.start()
thread.join()
print(messages) # ['hello']messages 没有从工作线程“传回”主线程。两条线程从一开始就在访问同一个 list 对象,join() 只保证主线程读取时工作线程已经结束。
共享本身不是错误。问题出现在多个线程同时依赖并修改一份可变状态,而且一次业务操作由多个步骤组成时。
一次“读取—计算—写回”如果被线程切换打断,可能发生这样的交错:
初始值 count = 0
Thread A 读取 count,得到 0
Thread B 读取 count,也得到 0
Thread A 写回 1
Thread B 写回 1
预期结果 2,实际结果 1最终结果依赖执行时序,这就是竞态条件(Race Condition)。Lock 可以把需要共同保护的读写操作包成临界区:
import threading
count = 0
lock = threading.Lock()
def increment() -> None:
global count
with lock:
current = count
count = current + 1with lock: 的执行过程是:
Thread A 尝试获取 lock
│
├── 获取成功
│ ├── 读取 count
│ ├── 计算新值
│ └── 写回 count
│
└── 离开 with,自动释放 lock
Thread B 在 A 持有 lock 时只能等待
│
└── A 释放后,B 才能进入同一临界区锁需要保护的是“读取—判断—修改”这一整组业务动作。只给写回那一行加锁仍然可能让两条线程基于同一个旧值计算。
锁也会带来等待,并且存在死锁风险。例如线程 A 持有 lock_a 等待 lock_b,同时线程 B 持有 lock_b 等待 lock_a,双方都无法继续。临界区保持短小、不同路径使用一致的加锁顺序,可以减少这种风险。
GIL 不能替代这里的业务锁。GIL 保护的是解释器内部状态,不保证一组业务操作具有原子性;线程也可能在阻塞 I/O、扩展模块执行或字节码调度期间交错。代码是否正确不能建立在“这次测试刚好没有出现竞态”上。
用 Queue 在线程之间传递数据
生产者和消费者直接共享 list 时,还要自行处理“列表为空时怎么等”“读写如何加锁”“任务何时结束”。标准库的 queue.Queue 已经提供线程安全的 put()、get() 和等待机制,适合作为线程之间的传送带。
生产者—消费者模型包含两个角色:
Producer Thread
│
│ 产生数据
▼
┌─────────────┐
│ Queue │
└─────────────┘
│
│ 取出数据
▼
Consumer Thread生产者不需要知道消费者正在执行到哪一行,消费者也不需要直接调用生产者。双方只约定数据格式,以及如何表示“后面没有数据了”。这比共享 list 少了一层相互依赖。
put() 把数据放到队尾,get() 从队头取走数据。它们对应的是数据传递,不是函数调用:
jobs.put("A")
item = jobs.get()当 Queue 中已经有数据时,get() 立即返回队头元素;Queue 为空时,调用它的线程进入阻塞状态:
Consumer 调用 jobs.get()
│
├── Queue 有数据 ──> 取出并返回
│
└── Queue 为空 ────> Consumer 暂停等待
│
Producer 调用 jobs.put("A") ─────┘
│
Consumer 被唤醒
│
get() 返回 "A"阻塞不等于 Python 在一个 while 循环里反复检查。等待期间,这条线程不持续执行 Python 代码,CPU 可以处理其他线程或其他进程的工作。
下面的例子把三个任务从生产者送到消费者,并完整处理启动、等待和退出:
import queue
import threading
import time
STOP = object()
jobs: queue.Queue[object] = queue.Queue()
def producer() -> None:
for item in ("A", "B", "C"):
jobs.put(item)
jobs.put(STOP)
def consumer() -> None:
while True:
item = jobs.get()
try:
if item is STOP:
return
time.sleep(0.2)
print(f"processed: {item}")
finally:
jobs.task_done()
producer_thread = threading.Thread(target=producer)
consumer_thread = threading.Thread(target=consumer)
producer_thread.start()
consumer_thread.start()
jobs.join()
producer_thread.join()
consumer_thread.join()Queue 为空时,get() 默认会阻塞消费者线程,直到有数据到达。此时线程处于等待状态,不会像下面的忙轮询一样持续占用 CPU:
# 不适合作为 Queue 的等待方式
while True:
if jobs.empty():
continueempty() 也不适合用来保证下一次 get() 一定成功。检查结束后、真正取数据前,其他消费者可能已经改变 Queue;这类“先检查再操作”的两步逻辑仍然存在竞态。直接使用阻塞 get(),或使用带超时的 get(timeout=...),才能把等待语义交给 Queue。
STOP 是哨兵值(Sentinel),表示后面没有更多数据。它不是线程的强制终止指令;消费者收到它后主动 return,目标函数结束,线程随之结束。使用多个消费者时,每个消费者都需要收到一个结束信号,否则仍有消费者会阻塞在 get()。
上面这段代码的 Queue 状态会这样变化:
| 动作 | Queue 内容 | 未完成任务计数 |
|---|---|---|
put("A") | A | 1 |
put("B") | A, B | 2 |
put("C") | A, B, C | 3 |
put(STOP) | A, B, C, STOP | 4 |
get() 取得 A | B, C, STOP | 仍为 4 |
完成 A 后 task_done() | B, C, STOP | 3 |
get() 只代表任务被取走,不代表处理已经完成。task_done() 才把未完成任务计数减一。哨兵也是一个通过 put() 入队的元素,所以消费者取到它后同样要调用一次 task_done();例子中的 finally 保证了这一点。
jobs.join() 与 thread.join() 等待的对象不同:前者等待所有入队项都被 task_done() 标记为处理完成,后者等待线程结束。相关计数规则见 queue.Queue 官方文档。
jobs.join() 等 Queue 中登记的工作全部处理完
thread.join() 等某一条线程的 target 函数结束两者经常同时出现,但不能互相替代。消费者忘记调用 task_done() 时,Queue 即使已经为空,jobs.join() 仍会一直等待;消费者的循环没有退出条件时,Queue 工作已经完成,thread.join() 仍会一直等待线程结束。
线程如何结束
目标函数正常 return 或抛出未处理异常时,线程就结束。threading.Thread 没有通用且安全的 stop() 方法,因为外部强行终止可能让锁、文件或事务停在未清理状态。
从代码视角看,一条线程会经历这些阶段:
创建 Thread 对象
│
│ 还没有执行 target
▼
调用 start()
│
│ target 开始执行
▼
运行 / 等待 I/O / 等待 Lock / 等待 Queue
│
│ target return 或抛出未处理异常
▼
线程结束线程在 sleep()、Queue.get() 或 Lock.acquire() 上等待时仍然是存活状态,只是暂时没有继续执行。is_alive() 回答的是线程是否从 start() 后尚未结束,不等于它此刻正在占用 CPU。
工作线程中的未处理异常默认不会沿着调用栈自动抛回主线程。Python 会调用 threading.excepthook() 输出异常信息,主线程仍可能继续。需要把结果和异常带回调用方时,ThreadPoolExecutor 返回的 Future 会更方便;直接使用 Thread 时通常要通过 Queue 或共享结果对象显式传递。
持续运行的工作线程通常通过协作信号退出。Queue 消费者可以使用哨兵值;周期性后台任务可以使用 Event:
import threading
stop_event = threading.Event()
def worker() -> None:
while not stop_event.wait(timeout=1):
print("heartbeat")
thread = threading.Thread(target=worker)
thread.start()
# 其他工作完成后发出停止信号
stop_event.set()
thread.join()Event 保存一个可以在线程之间共享的布尔信号。set() 把它置为真,is_set() 读取状态,wait(timeout) 会等待信号或超时。与单纯的 time.sleep(1) 相比,stop_event.wait(1) 可以在其他线程调用 set() 后提前醒来,退出响应更及时。
daemon=True 控制的是进程退出条件,不是正常关闭协议。当只剩 daemon 线程时,Python 可以退出进程,这些线程可能来不及释放文件、数据库事务等资源。需要清理资源的任务通常使用非 daemon 线程和明确的停止信号。官方文档也说明了 daemon 线程会在解释器关闭时被突然停止。
任务较多时使用线程池
任务数量增加后,逐个创建和保存 Thread 会混入较多调度代码。ThreadPoolExecutor 维护固定上限的工作线程,任务通过 submit() 进入线程池:
from concurrent.futures import ThreadPoolExecutor
def fetch(url: str) -> str:
return f"fetched: {url}"
urls = ["https://example.com/a", "https://example.com/b"]
with ThreadPoolExecutor(max_workers=5) as executor:
futures = [executor.submit(fetch, url) for url in urls]
results = [future.result() for future in futures]Future.result() 会等待对应任务并重新抛出任务中的异常。线程池限制并发线程数,却不能消除共享数据竞态;池中任务之间相互等待 Future 时还可能形成死锁。具体边界见 ThreadPoolExecutor 官方文档。
Thread 与线程池解决的层次不同:
| 形式 | 更接近什么 | 常见场景 |
|---|---|---|
Thread | 明确管理一条长期执行路径 | 后台监听、生产者线程、需要自己控制生命周期 |
ThreadPoolExecutor | 把许多独立任务交给一组复用线程 | 批量 HTTP 请求、批量文件 I/O |
线程也不是 Python 并发的唯一模型:
| 模型 | 执行方式 | 数据关系 | 主要边界 |
|---|---|---|---|
Thread | 多条操作系统线程 | 共享进程内存 | 共享状态需要同步;普通 CPython 受 GIL 影响 |
Process | 多个进程 | 默认内存隔离 | 进程创建和数据传递成本更高 |
asyncio | 通常在一条线程中调度多个 Task | 共享同一线程内对象 | 调用链需要配合 async / await,阻塞函数会卡住事件循环 |
普通 CPython 下,I/O 密集型同步代码常与 Thread 或线程池配合;大量纯 Python CPU 计算常转向进程;已经采用异步库的高并发 I/O 场景常使用 asyncio。这是执行模型的区别,不是固定的性能排名。
流式响应中的 Thread + Queue
长时间运行的生产任务与持续输出的 HTTP/SSE 响应可以放在两条执行路径中,用 Queue 传递每一段结果:
import queue
import threading
import time
from collections.abc import Iterator
STOP = object()
events: queue.Queue[object] = queue.Queue()
def long_running_stream() -> Iterator[str]:
for chunk in ("Hello", " ", "Thread"):
time.sleep(0.2)
yield chunk
def run_task() -> None:
try:
for chunk in long_running_stream():
events.put(chunk)
finally:
events.put(STOP)
def listen() -> Iterator[str]:
while True:
event = events.get()
if event is STOP:
break
yield str(event)
threading.Thread(target=run_task).start()
for piece in listen():
print(piece, end="", flush=True)Worker Thread HTTP / SSE Thread
│ │
long_running_stream() │
│ │
events.put(chunk) ───> Queue ───> events.get()
│
yield
│
Client这段代码可以按时间拆开:
1. 当前线程创建 Queue
2. 当前线程创建并启动 Worker Thread
3. Worker Thread 进入 long_running_stream()
4. 当前线程进入 listen(),调用 events.get()
5. Queue 为空,当前线程阻塞
6. Worker 产生 "Hello",调用 events.put("Hello")
7. 当前线程被唤醒,get() 返回 "Hello"
8. listen() yield "Hello"
9. 当前线程再次调用 get(),等待下一段数据
10. Worker 结束,在 finally 中放入 STOP
11. listen() 取得 STOP,break,生成器结束核心不是 yield 自己创建了线程,而是两条执行路径承担了不同工作:工作线程可以长时间生成数据,当前线程可以一边等待 Queue,一边把已经到达的数据交给调用方。
如果所有工作都放在同一条执行路径中,通常会变成:
先完整运行 long_running_stream()
│
└── 全部结束后才开始消费和返回拆开后则是:
生成一段 ──> Queue ──> 返回一段
生成一段 ──> Queue ──> 返回一段
生成结束 ──> STOP ──> 响应结束接入 Web 框架后,listen() 可以作为 HTTP/SSE 响应体的生成器。这里假设的是同步服务器中的请求执行路径;具体由线程、进程还是其他并发单元承载,要看应用服务器配置。
finally 中的结束信号只能保证消费者有机会退出等待,不能把工作线程的异常详情送给客户端。如果流式任务的错误也需要展示,Queue 中通常传递带类型的事件,例如 {"type": "data", "value": chunk}、{"type": "error", "message": ...} 和 {"type": "done"}。实际服务还要处理客户端断开、超时、Queue 容量、资源清理和应用服务器的线程模型;这些属于建立基本数据流之后的下一层问题。