Python 原生支持多线程编程,标准库内置 threading、_thread 模块,开箱即用。
但是,受 GIL(全局解释器锁) 限制,CPU 密集型任务无法利用多核,IO 密集型任务多线程提升巨大。
要想了解线程,我们必须先知道什么是进程?
进程是操作系统分配资源的最小单位,每个程序独立占内存、CPU 资源。比如打开 2 个 Python 脚本就是 2 个进程,互相隔离。又比如打开音乐软件听歌,也是一个进程。打开聊天软件也是一个进程。
那么,什么是线程?
线程是进程内部最小执行单元,一个进程可以包含多条线程,共享该进程的内存、文件、网络资源。一个进程默认有 1 个主线程。多线程让一个程序同时执行多个任务,共享同一块内存空间,切换成本低。
如果还觉得不好理解,可以这么想?进程 = 一间工厂、线程 = 工厂里的工人。同一工厂工人共用水电设备,不同工厂完全隔。
知道了线程,多线程又是什么?
说到多线程,不得不提一下单线程,单线程指代码从上到下顺序执行,一件事做完再做下一件,遇到等待(网络、文件、浏览器加载)全程卡死。根据上面工厂的比喻,即一个工厂就一个工人。
那多线程呢!指在同一个程序里,同时启动多条执行支路,并发处理多个任务,等待期间切换其他任务干活。
多线程的作用:
(1)提高程序并发效率,如爬虫、文件下载、UI 界面不卡顿。
(2)处理 I/O 密集型任务,如网络请求、文件读写非常高效。
注意:
Python 有 GIL(全局解释器锁),多线程不适合 CPU 密集型任务,但非常适合 I/O 密集型任务。
GIL(Global Interpreter Lock,全局解释器锁) 是 CPython(我们日常使用的标准 Python 解释器)内部的一把互斥锁。同一时刻,一个 Python 进程内,永远只能有 1 个线程执行真正的 Python 字节码。
通过一个简单比喻说明:一间房子(Python 进程)只有一把钥匙(GIL 锁),多个线程(人)想干活,必须先拿到钥匙。同一时间钥匙只能给一个人,其他人只能排队等钥匙。
Python 提供两个多线程模块:
_thread:低级模块,简单但功能少(不推荐生产使用)
threading:高级模块,功能完善,官方推荐
下面将分别介绍这两个模块。
_thread 是 Python 内置的轻量线程模块,语法简单,适合快速入门。
基本用法如下:
# 1.导入模块
import _thread
import time
# 2.定义线程执行的函数
def task(thread_name, delay):
for i in range(3):
time.sleep(delay)
print(f"线程 {thread_name} 执行:第{i+1}次")
# 3.创建并启动线程
# 参数:(函数名, (参数元组,))
_thread.start_new_thread(task, ("线程-1", 1))
_thread.start_new_thread(task, ("线程-2", 1.5))
# 4.主线程等待
# 必须添加,如果不添加,Python 脚本直接就执行完了,主线程结束,子线程也就跟着结束了
time.sleep(5)
print("主线程结束")使用 _thread 实现多线程非常简单,只需导入依赖,通过 _thread.start_new_thread() 函数将另一个函数作为参数,就可以快速创建一个新的线程。函数的参数可以通过 start_new_thread() 的第二个参数进行传递。
运行示例,输出如下:
线程 线程-1 执行:第1次
线程 线程-2 执行:第1次
线程 线程-1 执行:第2次
线程 线程-2 执行:第2次
线程 线程-1 执行:第3次
线程 线程-2 执行:第3次
主线程结束注意,虽然用法非常简单,但是缺陷也很明显,使用 start_new_thread() 创建的线程没有线程等待、没有线程名称、没有守护线程,生产环境不推荐使用。
threading 是 Python 官方标准库,功能完善,生产环境推荐使用这个模块。
主要支持如下功能:
线程等待(join)
守护线程(daemon)
线程名称、当前线程获取
继承方式创建线程
下面将分别介绍使用 threading 模块创建线程的几种方式:
该种方式和 _thread 模块的 start_new_thread() 方式一样,均是将一个函数传递给模块提供的函数,来快速创建线程,如下:
import threading
import time
# 线程函数
def task(name):
print(f"线程 {name} 开始")
time.sleep(2)
print(f"线程 {name} 结束")
# 创建线程
t1 = threading.Thread(target=task, args=("A",))
t2 = threading.Thread(target=task, args=("B",))
# 启动线程
t1.start()
t2.start()
# 主线程等待子线程完成
t1.join()
t2.join()
print("所有线程执行完毕")运行示例,输出如下:
线程 A 开始
线程 B 开始
线程 B 结束线程 A 结束
所有线程执行完毕通过继承 Thread 类实现创建线程,例如:
import threading
import time
# 自定义线程类
class MyThread(threading.Thread):
def __init__(self, name):
super().__init__()
self.name = name
# 线程必须重写 run 方法
def run(self):
print(f"线程 {self.name} 运行")
time.sleep(2)
print(f"线程 {self.name} 结束")
# 创建并启动
t1 = MyThread("线程1")
t2 = MyThread("线程2")
t1.start()
t2.start()上面代码没有手动调用两个线程的 join() 方法,这是为什么?注意,不是不需要 join(),只是上面的代码不用等待,程序依然会等线程跑完才退出,这是 Python 主线程的守护线程机制在起作用。
线程常用方法如下:
threading.current_thread() 获取当前线程
threading.active_count() 获取活跃线程数
t1.start() 启动线程
t1.name 设置线程名/获取线程名
t1.is_alive() 判断线程是否存活
t1.daemon 设置守护线程(主线程退出,子线程自动退出)
简单示例,演示上面方法的用法:
import threading
import time
def work_task():
# 使用 threading.current_thread() 获取当前执行的线程对象
current_t = threading.current_thread()
print(f"【任务内】当前线程对象:{current_t}")
print(f"【任务内】当前线程名称:{current_t.name}\n")
# 模拟任务执行3秒
for i in range(3):
print(f"[{current_t.name}] 正在工作,进度 {i+1}/3")
time.sleep(1)
if __name__ == "__main__":
# 主线程信息
main_thread = threading.current_thread()
print("===== 主线程初始信息 =====")
print(f"主线程对象:{main_thread}")
print(f"主线程名称:{main_thread.name}")
print(f"当前活跃线程总数:{threading.active_count()}\n")
# 创建子线程
t1 = threading.Thread(target=work_task)
# 设置线程名
t1.name = "测试线程"
# 获取线程名
print(f"子线程设置后的名称:{t1.name}\n")
# 设置为守护线程
# 守护线程规则:主线程代码执行完毕,直接杀死守护子线程,不等它跑完
t1.daemon = True
# 启动线程
t1.start()
print("===== 启动子线程后 =====")
# 获取当前存活的线程总数(主线程 + t1)
print(f"启动后活跃线程数量:{threading.active_count()}")
# 判断线程是否还在运行
print(f"子线程是否存活:{t1.is_alive()}\n")
# 主线程只等待2秒就结束
print("主线程等待2秒后退出...")
time.sleep(2)
print("\n主线程执行完毕,准备退出程序!")
# 因为t1是守护线程,主线程结束后,子线程会被强制终止,不会完整跑完3秒
print(f"主线程退出前,子线程是否还存活:{t1.is_alive()}")运行示例,输出如下:
===== 主线程初始信息 =====
主线程对象:<_MainThread(MainThread, started 16932)>
主线程名称:MainThread
当前活跃线程总数:1
子线程设置后的名称:测试线程
【任务内】当前线程对象:<Thread(测试线程, started daemon 8832)>
===== 启动子线程后 =====
【任务内】当前线程名称:测试线程
启动后活跃线程数量:2
[测试线程] 正在工作,进度 1/3子线程是否存活:True
主线程等待2秒后退出...
[测试线程] 正在工作,进度 2/3
主线程执行完毕,准备退出程序!
主线程退出前,子线程是否还存活:True假如有多个线程同时修改全局变量会发生什么?数据会错乱,为什么会错乱,下面分析一下:
假如存在一个计数的 count 全局变量,每个线程通过 count = count + 1 实现计数功能。然而, count = count + 1 并不是一个原则操作,会被拆分成三步,三步执行过程中线程可以被切换:
步骤 1:读取内存中全局变量 count 的值(LOAD_GLOBAL)。注意:这里可能被多个线程同时获取到该值。
步骤 2:拿到的值 + 1 做加法运算(INCREMENT)。此时另一个线程已经完成了+1 操作。而当前线程依然用步骤 1 获取的旧值进行+1 操作,最终会导致已完成线程的+1 操作丢失。
步骤 3:把计算结果写回全局变量 count(STORE_GLOBAL)
多线程既然存在并发问题,该如何解决呢?为了解决该问题,需要引入线程锁(Lock)。Python 中,可以通过 threading.Lock() 来保证同一时间只有一个线程修改数据。通过线程锁实现 count = count + 1 操作同一时刻只在一个线程中执行。
注意,在 python 中,count += 1 操作不是原子的,例如:
import dis
count = 0
def add():
global count
count += 1
dis.dis(add)代码输出如下:
4 RESUME 0
6 LOAD_GLOBAL 0 (count) # 读取全局count
LOAD_CONST 1 (1) # 加载常量1
BINARY_OP 13 (+=) # 加法运算
STORE_GLOBAL 0 (count) # 写回count
RETURN_CONST 0 (None)(1)如果没有对多线程中访问的共享资源添加锁,会出现意想不到的问题。如全局变量 count,使用多个线程对 count 进行计数,如下:
import threading
import hashlib
import time
count = 0
def add():
global count
for i in range(100000):
# 为了更容易触发 count += 1 的多线程问题
# 特意将 count += 1 拆分为三部,并且使用 md5 加密消耗CPU时间
temp = count
get_str_md5(str(i)) # 模拟耗时操作,触发线程切换
temp += 1
get_str_md5(str(i)) # 模拟耗时操作,触发线程切换
count = temp
# 使用MD5模拟一个耗时操作
def get_str_md5(content: str, encoding="utf-8") -> str:
md5_obj = hashlib.md5()
md5_obj.update(content.encode(encoding))
md5_result = md5_obj.hexdigest()
return md5_result
# 测试多线程下的全局变量数据错乱问题
if __name__ == "__main__":
exeCount = 0
while True:
count = 0
exeCount += 1
t1 = threading.Thread(target=add)
t2 = threading.Thread(target=add)
t1.start()
t2.start()
t1.join()
t2.join()
print("执行次数:", exeCount)
if(count % 2 != 0):
print("数据错乱,最终结果:", count)
break运行上面代码,如果没有出现数据错乱情况,可以添加 md5 执行的个数,或者多执行几次就会出现。我也是折腾了很久才出现数据错乱的现象。
运行示例,输出如下:
执行次数: 1
执行次数: 2
数据错乱,最终结果: 112447使用 threading 的 Lock() 创建锁,然后通过 acquire() 加锁,release() 释放锁。例如:
import threading
import hashlib
import time
count = 0
lock = threading.Lock() # 创建锁
def add():
global count
for i in range(100000):
try:
lock.acquire() # 加锁
# 为了更容易触发 count += 1 的多线程问题,特意将 count += 1 拆分为三部,并且使用 md5 加密消耗CPU时间
temp = count
get_str_md5(str(i)) # 模拟耗时操作,触发线程切换
temp += 1
get_str_md5(str(i)) # 模拟耗时操作,触发线程切换
count = temp
finally:
lock.release() # 释放锁
# 使用MD5模拟一个耗时操作
def get_str_md5(content: str, encoding="utf-8") -> str:
md5_obj = hashlib.md5()
md5_obj.update(content.encode(encoding))
md5_result = md5_obj.hexdigest()
return md5_result
# 测试多线程下的全局变量数据错乱问题
if __name__ == "__main__":
exeCount = 0
while True:
count = 0
exeCount += 1
t1 = threading.Thread(target=add)
t2 = threading.Thread(target=add)
t1.start()
t2.start()
t1.join()
t2.join()
print("执行次数:", exeCount)
if(count % 2 != 0):
print("数据错乱,最终结果:", count)
break运行代码,数据错乱问题就解决了。但是,加锁太繁琐了,有没有更简单的方法呢!有的,请使用 with 自动锁机制,例如,将加锁部分的代码改为:
def add():
global count
for i in range(100000):
with lock: # 使用锁来保护对 count 的访问
# 为了更容易触发 count += 1 的多线程问题
# 特意将 count += 1 拆分为三部,并且使用 md5 加密消耗CPU时间
temp = count
get_str_md5(str(i)) # 模拟耗时操作,触发线程切换
temp += 1
get_str_md5(str(i)) # 模拟耗时操作,触发线程切换
count = temp上面使用 with lock 是不是非常简单,加锁和释放锁由 Python 完成。
queue 是 Python 线程安全的队列模块,专门用于多线程数据通信,自动处理锁机制。
提供三种队列:
Queue:FIFO 先进先出队列
LifoQueue:LIFO 后进先出队列(栈)
PriorityQueue:优先级队列(数字越小优先级越高)
queue.Queue 是 Python 标准库提供的线程安全先进先出 (FIFO) 队列,专门用于多线程之间安全传递数据,自带锁机制,天然解决多线程读写竞态问题,无需手动 threading.Lock。
先进先出:先放进去的数据,先取出来。
内置同步锁,多线程并发 put /get 不会数据错乱。
支持阻塞、非阻塞、队列容量限制。
常用的方法介绍:
Queue(maxsize=0) 初始化队列,maxsize=0 代表无容量上限
put(item, block=True, timeout=None) 放入数据,队列满时阻塞等待
get(block=True, timeout=None) 取出数据;队列为空时阻塞等待
qsize() 返回当前队列元素个数(近似值,多线程下不准)
empty() 判断队列是否为空(多线程下有瞬时误差)
full() 判断队列是否已满
task_done() 任务完成标记,配合 join () 使用
join() 阻塞主线程,直到队列所有任务都调用 task_done ()
详细信息参考 https://docs.python.org/zh-cn/3.11/library/queue.html 手册。
示例:简单使用 queue 队列。
import queue
# 创建队列
q = queue.Queue()
# 放数据
q.put("A")
q.put("B")
q.put("C")
# 取数据
print(q.get()) # A
print(q.get()) # B
print(q.get()) # C示例:在多线程生产消费模型中使用 queue 来进行通信。
import queue
import threading
import time
import random
# 创建一个队列,最大容量为5
q = queue.Queue(5)
# 生产者
def producer():
for i in range(5):
q.put(f"产品{i}")
print(f"生产:产品{i}")
time.sleep(random.uniform(0.1, 0.5)) # 模拟生产时间
# 生产完毕,放入哨兵,通知消费者结束
q.put(None)
# 消费者
def consumer():
while True:
data = q.get()
# 读到哨兵,退出循环
if data is None:
q.task_done()
break
print(f"消费:{data}")
time.sleep(random.uniform(0.1, 0.5)) # 模拟消费时间
q.task_done()
t1 = threading.Thread(target=producer)
# 设置为守护线程,主线程结束时会自动结束
t2 = threading.Thread(target=consumer, daemon=True)
t1.start()
t2.start()
t1.join()
t2.join()
print("生产者和消费者线程已结束。")上述示例代码,生产者将生产 5 个产品,当产品都生产完成后,放入一个特殊值 None 到队列。当消费者从队列消费消息时,消费到 None,则自动结束,退出 while 循环。
运行示例,效果如下:
生产:产品0
消费:产品0
生产:产品1
生产:产品2
消费:产品1
消费:产品2
生产:产品3
生产:产品4
消费:产品3
消费:产品4
生产者和消费者线程已结束。queue.LifoQueue 是 Python 标准库 queue 提供的后进先出队列,逻辑等价于「栈 Stack」:
最后放进去的数据,最先取出来;
和 Queue(FIFO) 底层共享同一套线程安全锁、阻塞机制、API;
多线程并发读写天然安全,无需手动加锁;
示例:
from queue import LifoQueue
# 无容量限制栈
stack = LifoQueue()
# 入栈(压栈)
stack.put("A")
stack.put("B")
stack.put("C")
# 后进先出,先取出最后放入的C
print(stack.get()) # C
print(stack.get()) # B
print(stack.get()) # Aqueue.PriorityQueue 是标准库提供的线程安全优先队列,底层基于堆(heapq)实现:
出队规则:最小值优先弹出(数字越小,优先级越高,越先被取出);
线程安全:内部自带锁,多线程 put/get 无数据竞争;
接口、阻塞机制、task_done()/join() 与 Queue、LifoQueue 完全统一;
示例:简单了解优先级队列的用法。
from queue import PriorityQueue
pq = PriorityQueue()
# 存入 (优先级, 数据)
pq.put((3, "普通任务"))
pq.put((1, "紧急任务"))
pq.put((2, "一般任务"))
# 按优先级从小到大取出,1 > 2 > 3
while not pq.empty():
prio, task = pq.get()
print(f"优先级{prio}:{task}")
pq.task_done()运行示例,输出如下:
优先级1:紧急任务
优先级2:一般任务
优先级3:普通任务示例 2:利用优先级队列实现任务按照优先级优先处理,当生产者向队列中放入一个高优先级的结束标记时,消费者停止处理后续的低优先级任务。
import queue
import threading
import time
import random
# 最大容量6
pq = queue.PriorityQueue(maxsize=6)
def producer():
# 生成5个不同优先级任务
for i in range(5):
prio = random.randint(1, 5)
task_name = f"任务{i}"
pq.put((prio, task_name))
print(f"存入:优先级{prio} {task_name}")
time.sleep(random.uniform(0.1, 0.3))
# 哨兵结束信号
pq.put((0, None)) # 优先级很高,会导致低优先级任务不被处理
def consumer():
while True:
prio, data = pq.get()
# 哨兵消息优先级很高,会优先被执行
if data is None:
print("收到结束信号,退出")
pq.task_done()
break
print(f"【处理】优先级{prio} → {data}")
time.sleep(random.uniform(0.2, 0.5))
pq.task_done()
t_pro = threading.Thread(target=producer)
t_con = threading.Thread(target=consumer)
t_pro.start()
t_con.start()
t_pro.join()
t_con.join()
print("全部任务处理完毕")运行实例,输出如下:
存入:优先级3 任务0
【处理】优先级3 → 任务0
存入:优先级4 任务1
【处理】优先级4 → 任务1
存入:优先级5 任务2
存入:优先级1 任务3
【处理】优先级1 → 任务3
存入:优先级3 任务4
【处理】优先级3 → 任务4
收到结束信号,退出
全部任务处理完毕在前面使用的 put() 方法时,仅仅传递了值,get() 方法一个参数也没有指定。其实,这两个方法还支持如下参数:
block 用来控制队列满时是否阻塞等待,True 表示阻塞等待,False 表示不阻塞等待。
timeout 用来控制阻塞等待的时长,None 表示一直等待。
示例:演示不阻塞和等待超时效果。
from queue import PriorityQueue, Full, Empty
# 创建一个优先级队列,最大存放1个元素
pq = PriorityQueue(maxsize=1)
pq.put((1, "测试"))
# 非阻塞 put,队列满抛 Full 异常
try:
pq.put((2, "新任务"), block=False) # block 为 False 不阻塞
except Full:
print("队列已满")
# 从队列取一个值
print(pq.get()) # (1, '测试')
# 此时的队列中为空,没有值了
# 超时get,2秒取不到抛 Empty 异常
try:
pq.get(timeout=2)
except Empty:
print("2秒无数据")运行示例,输出如下:
队列已满
(1, '测试')
2秒无数据队列内部维护一个未完成任务计数器 unfinished_tasks:
调用 q.put(item) → unfinished_tasks += 1
消费完任务调用 q.task_done() → unfinished_tasks -= 1
q.join() 会阻塞当前线程,直到 unfinished_tasks == 0
只有所有入队任务都被标记完成,join() 才会放行,主线程继续向下执行
示例:演示使用队列的 join() 等待队列所有任务被工作线程处理完成。
from queue import PriorityQueue
import threading
# 创建优先级队列
pq = PriorityQueue()
# 定义工作线程函数
def worker():
while True:
prio, data = pq.get()
print("处理...", data)
pq.task_done() # 必须添加,表示任务完成
t = threading.Thread(target=worker, daemon=True) # 后台运行
t.start()
# 批量放入任务
pq.put((5, "C"))
pq.put((1, "A"))
pq.put((3, "B"))
# 等待所有任务处理完成
pq.join()
print("所有优先级任务执行结束")运行示例,输出如下:
处理... A
处理... B
处理... C
所有优先级任务执行结束到这里,Python 线程基础知识就介绍完了。后续将介绍多线程中的可重入锁、信号量等相关知识。