🚀 25. 并发编程

本章概述

学习 Python 并发编程:多线程、多进程、协程,提升程序效率。 预计学习时间:60 分钟


25.1 并发 vs 并行

概念

  • 并发(Concurrency):多个任务交替执行,看起来像同时进行
  • 并行(Parallelism):多个任务真正同时执行(需要多核CPU)

Python 的并发方式

方式适用场景特点
多线程I/O 密集型任务有GIL限制,不能真正并行
多进程CPU 密集型任务可以真正并行,开销大
协程I/O 密集型任务单线程内切换,开销最小

GIL(全局解释器锁)

Python 有 GIL(Global Interpreter Lock),同一时刻只能有一个线程执行 Python 字节码。 所以多线程不能真正并行执行 CPU 密集型任务,但 I/O 密集型任务可以受益。


25.2 多线程 threading

基本用法

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("所有任务完成")

继承 Thread 类

import threading
import time
 
class MyThread(threading.Thread):
    def __init__(self, name):
        super().__init__()
        self.name = name
 
    def run(self):  # 线程执行的代码
        print(f"线程 {self.name} 开始")
        time.sleep(2)
        print(f"线程 {self.name} 结束")
 
t1 = MyThread("A")
t2 = MyThread("B")
 
t1.start()
t2.start()
 
t1.join()
t2.join()

线程安全问题

多个线程同时修改共享数据,会出现问题。

import threading
 
count = 0
 
def increment():
    global count
    for _ in range(100000):
        count += 1
 
t1 = threading.Thread(target=increment)
t2 = threading.Thread(target=increment)
 
t1.start()
t2.start()
t1.join()
t2.join()
 
print(count)  # 不是200000!因为线程不安全

锁 Lock

import threading
 
count = 0
lock = threading.Lock()
 
def increment():
    global count
    for _ in range(100000):
        lock.acquire()  # 获取锁
        count += 1
        lock.release()  # 释放锁
 
# 或者用 with 语句(推荐)
def increment():
    global count
    for _ in range(100000):
        with lock:  # 自动获取和释放
            count += 1
 
t1 = threading.Thread(target=increment)
t2 = threading.Thread(target=increment)
 
t1.start()
t2.start()
t1.join()
t2.join()
 
print(count)  # 200000(正确了)

常用方法

# 当前线程
threading.current_thread()
 
# 活跃线程数
threading.active_count()
 
# 所有线程列表
threading.enumerate()
 
# 线程是否存活
t.is_alive()
 
# 设置守护线程(主线程结束时自动结束)
t = threading.Thread(target=task, daemon=True)

25.3 多进程 multiprocessing

基本用法

import multiprocessing
import time
 
def task(name):
    print(f"进程 {name} 开始")
    time.sleep(2)
    print(f"进程 {name} 结束")
 
if __name__ == "__main__":
    # 创建进程
    p1 = multiprocessing.Process(target=task, args=("A",))
    p2 = multiprocessing.Process(target=task, args=("B",))
 
    # 启动进程
    p1.start()
    p2.start()
 
    # 等待进程结束
    p1.join()
    p2.join()
 
    print("所有进程完成")

注意

Windows 下多进程代码必须放在 if __name__ == "__main__": 里面,否则会报错。

进程池 Pool

import multiprocessing
import time
 
def task(n):
    time.sleep(1)
    return n * n
 
if __name__ == "__main__":
    # 创建进程池,大小为CPU核心数
    with multiprocessing.Pool() as pool:
        # map:批量执行
        results = pool.map(task, [1, 2, 3, 4, 5])
        print(results)  # [1, 4, 9, 16, 25]
 
        # apply_async:异步执行
        result = pool.apply_async(task, (10,))
        print(result.get())  # 100

进程间通信

Queue 队列

import multiprocessing
 
def producer(q):
    for i in range(5):
        q.put(i)
        print(f"生产:{i}")
 
def consumer(q):
    while True:
        item = q.get()
        if item is None:
            break
        print(f"消费:{item}")
 
if __name__ == "__main__":
    q = multiprocessing.Queue()
 
    p1 = multiprocessing.Process(target=producer, args=(q,))
    p2 = multiprocessing.Process(target=consumer, args=(q,))
 
    p1.start()
    p2.start()
 
    p1.join()
    q.put(None)  # 发送结束信号
    p2.join()

25.4 协程 asyncio

协程是比线程更轻量级的并发方式,在单线程内切换。

基本用法

import asyncio
 
async def task(name, seconds):
    print(f"任务 {name} 开始")
    await asyncio.sleep(seconds)  # 异步等待
    print(f"任务 {name} 结束")
    return f"{name}完成"
 
async def main():
    # 方式1:顺序执行(慢)
    # await task("A", 2)
    # await task("B", 1)
 
    # 方式2:并发执行(快)
    task1 = task("A", 2)
    task2 = task("B", 1)
 
    results = await asyncio.gather(task1, task2)
    print(results)
 
# 运行协程
asyncio.run(main())

创建任务

import asyncio
 
async def task(name, seconds):
    print(f"任务 {name} 开始")
    await asyncio.sleep(seconds)
    print(f"任务 {name} 结束")
    return name
 
async def main():
    # 创建任务
    t1 = asyncio.create_task(task("A", 2))
    t2 = asyncio.create_task(task("B", 1))
 
    # 等待任务完成
    result1 = await t1
    result2 = await t2
    print(result1, result2)
 
asyncio.run(main())

asyncio.gather

并发运行多个协程:

import asyncio
 
async def task(name, seconds):
    await asyncio.sleep(seconds)
    return name
 
async def main():
    results = await asyncio.gather(
        task("A", 2),
        task("B", 1),
        task("C", 3)
    )
    print(results)  # ['A', 'B', 'C']
 
asyncio.run(main())

25.5 并发编程对比

什么时候用什么?

场景推荐方式原因
I/O 密集型(网络请求、文件读写)协程 / 多线程等待时可以切换,效率高
CPU 密集型(计算、数据处理)多进程绕过GIL,真正并行
简单任务,代码简单多线程简单易用
高并发 I/O 任务协程开销最小,效率最高

对比总结

维度多线程多进程协程
切换开销
数据共享容易(注意锁)困难(需要IPC)容易
并行能力不能(GIL)不能
适用场景I/O密集型CPU密集型I/O密集型
代码复杂度

25.6 常用并发库

concurrent.futures

统一的线程池和进程池接口:

from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
import time
 
def task(n):
    time.sleep(1)
    return n * n
 
# 线程池
with ThreadPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(task, [1, 2, 3, 4, 5]))
    print(results)
 
# 进程池
with ProcessPoolExecutor(max_workers=4) as executor:
    results = list(executor.map(task, [1, 2, 3, 4, 5]))
    print(results)

submit 和 as_completed

from concurrent.futures import ThreadPoolExecutor, as_completed
import time
 
def task(name, seconds):
    time.sleep(seconds)
    return f"{name}完成"
 
with ThreadPoolExecutor() as executor:
    futures = [
        executor.submit(task, "A", 3),
        executor.submit(task, "B", 1),
        executor.submit(task, "C", 2)
    ]
 
    # 按完成顺序获取结果
    for future in as_completed(futures):
        print(future.result())

🔗 相关章节


📝 我的笔记

在这里记录你的理解和练习代码

# 你的练习代码
 

✅ 本章检查清单

  • 理解并发和并行的概念
  • 了解 GIL 的影响
  • 会使用多线程 threading
  • 了解线程安全和锁
  • 会使用多进程 multiprocessing
  • 了解进程池的使用
  • 了解协程 asyncio 的基本用法
  • 知道不同并发方式的适用场景
  • 了解 concurrent.futures