Python3 多线程编程

🎉摘要:本文全面讲解Python多线程编程,涵盖_thread和threading模块的使用,深入解析GIL(全局解释器锁)对多线程的影响。详细介绍线程创建、同步、线程安全队列(Queue/LifoQueue/PriorityQueue)等核心知识,并提供大量实战代码示例,帮助开发者高效处理IO密集型任务。

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 多线程编程

Python 提供两个多线程模块:

  • _thread:低级模块,简单但功能少(不推荐生产使用)

  • threading:高级模块,功能完善,官方推荐

下面将分别介绍这两个模块。

_thread 模块(低级模块)

_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 模块(高级模块)

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 类创建线程(面向对象)

通过继承 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 模块(线程安全队列)

queue 是 Python 线程安全的队列模块,专门用于多线程数据通信,自动处理锁机制。

提供三种队列:

  1. Queue:FIFO 先进先出队列

  2. LifoQueue:LIFO 后进先出队列(栈)

  3. PriorityQueue:优先级队列(数字越小优先级越高)

FIFO 队列 Queue(最常用)

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
生产者和消费者线程已结束。

LIFO 队列 LifoQueue

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())  # A

优先级队列 PriorityQueue

queue.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秒无数据

task_done () + join () 配合使用

队列内部维护一个未完成任务计数器 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 线程基础知识就介绍完了。后续将介绍多线程中的可重入锁、信号量等相关知识。

说说我的看法
全部评论(
没有评论
关于
本网站专注于 Java、数据库(MySQL、Oracle)、Linux、软件架构及大数据等多领域技术知识分享。涵盖丰富的原创与精选技术文章,助力技术传播与交流。无论是技术新手渴望入门,还是资深开发者寻求进阶,这里都能为您提供深度见解与实用经验,让复杂编码变得轻松易懂,携手共赴技术提升新高度。如有侵权,请来信告知:hxstrive@outlook.com
其他应用
公众号