Python多线程与多进程:并发编程实战

文章总体概览信息图

引言

并发编程是现代后端开发的核心技能,它能够显著提升应用程序的性能和响应能力。作为一名从Python转向Rust的后端开发者,我在实践中总结了Python并发编程的最佳实践。本文将深入探讨Python的多线程和多进程编程,帮助你构建高效的并发系统。

一、并发编程基础概念

1.1 什么是并发

并发是指程序在同一时间段内处理多个任务的能力。

1.2 并发与并行的区别

概念 说明
并发 多个任务交替执行,宏观上看起来同时进行
并行 多个任务真正同时执行,需要多个CPU核心

1.3 Python中的并发方式

  • 多线程:适合IO密集型任务
  • 多进程:适合CPU密集型任务
  • 协程:轻量级并发,适合高IO场景

二、多线程编程

2.1 threading模块

import threading
import time

def worker(num):
    print(f"Worker {num} started")
    time.sleep(2)
    print(f"Worker {num} finished")

threads = []
for i in range(5):
    t = threading.Thread(target=worker, args=(i,))
    threads.append(t)
    t.start()

for t in threads:
    t.join()

2.2 线程同步

import threading

counter = 0
lock = threading.Lock()

def increment():
    global counter
    for _ in range(100000):
        with lock:
            counter += 1

threads = []
for _ in range(10):
    t = threading.Thread(target=increment)
    threads.append(t)
    t.start()

for t in threads:
    t.join()

print(f"Final counter: {counter}")

2.3 线程局部存储

import threading

local_data = threading.local()

def worker():
    local_data.value = threading.current_thread().name
    print(f"{threading.current_thread().name}: {local_data.value}")

threads = []
for i in range(3):
    t = threading.Thread(target=worker, name=f"Thread-{i}")
    threads.append(t)
    t.start()

for t in threads:
    t.join()

2.4 线程池

from concurrent.futures import ThreadPoolExecutor
import time

def task(name):
    print(f"Task {name} started")
    time.sleep(2)
    print(f"Task {name} finished")
    return f"Result {name}"

with ThreadPoolExecutor(max_workers=3) as executor:
    futures = [executor.submit(task, i) for i in range(5)]
    
    for future in futures:
        result = future.result()
        print(f"Received: {result}")

三、多进程编程

3.1 multiprocessing模块

import multiprocessing
import time

def worker(num):
    print(f"Worker {num} started")
    time.sleep(2)
    print(f"Worker {num} finished")

if __name__ == "__main__":
    processes = []
    for i in range(5):
        p = multiprocessing.Process(target=worker, args=(i,))
        processes.append(p)
        p.start()

    for p in processes:
        p.join()

3.2 进程间通信

import multiprocessing

def producer(queue):
    for i in range(5):
        queue.put(i)
        print(f"Produced: {i}")

def consumer(queue):
    while True:
        item = queue.get()
        if item is None:
            break
        print(f"Consumed: {item}")
        queue.task_done()

if __name__ == "__main__":
    queue = multiprocessing.JoinableQueue()
    
    p1 = multiprocessing.Process(target=producer, args=(queue,))
    p2 = multiprocessing.Process(target=consumer, args=(queue,))
    
    p1.start()
    p2.start()
    
    p1.join()
    queue.put(None)
    p2.join()

3.3 进程池

from concurrent.futures import ProcessPoolExecutor
import time

def compute_square(n):
    time.sleep(1)
    return n * n

if __name__ == "__main__":
    with ProcessPoolExecutor(max_workers=4) as executor:
        numbers = [1, 2, 3, 4, 5]
        results = executor.map(compute_square, numbers)
        
        for num, result in zip(numbers, results):
            print(f"{num}^2 = {result}")

3.4 共享内存

import multiprocessing

def update_counter(counter):
    for _ in range(100000):
        counter.value += 1

if __name__ == "__main__":
    counter = multiprocessing.Value('i', 0)
    
    processes = []
    for _ in range(10):
        p = multiprocessing.Process(target=update_counter, args=(counter,))
        processes.append(p)
        p.start()

    for p in processes:
        p.join()

    print(f"Final counter: {counter.value}")

四、线程与进程的对比

4.1 性能对比

特性 线程 进程
内存开销
启动速度
CPU利用率 受GIL限制 充分利用多核
数据共享 方便 需要IPC
适用场景 IO密集型 CPU密集型

4.2 GIL的影响

Python的全局解释器锁(GIL)限制了多线程程序的CPU利用率。对于CPU密集型任务,多进程是更好的选择。

五、实用案例

5.1 爬取多个网页

import requests
from concurrent.futures import ThreadPoolExecutor

def fetch_url(url):
    response = requests.get(url)
    return url, response.status_code

urls = [
    "https://example.com",
    "https://google.com",
    "https://github.com",
    "https://python.org",
    "https://rust-lang.org"
]

with ThreadPoolExecutor(max_workers=5) as executor:
    results = executor.map(fetch_url, urls)
    
    for url, status in results:
        print(f"{url}: {status}")

5.2 并行数据处理

import numpy as np
from concurrent.futures import ProcessPoolExecutor

def process_chunk(chunk):
    return chunk.sum()

if __name__ == "__main__":
    data = np.random.rand(10_000_000)
    
    chunks = np.array_split(data, 4)
    
    with ProcessPoolExecutor(max_workers=4) as executor:
        results = executor.map(process_chunk, chunks)
        
        total = sum(results)
        print(f"Total sum: {total}")

5.3 异步任务队列

import time
from concurrent.futures import ThreadPoolExecutor

class TaskQueue:
    def __init__(self, max_workers=5):
        self.executor = ThreadPoolExecutor(max_workers=max_workers)
        self.futures = []
    
    def submit(self, func, *args):
        future = self.executor.submit(func, *args)
        self.futures.append(future)
    
    def wait(self):
        for future in self.futures:
            future.result()
        self.executor.shutdown()

def process_task(task_id):
    print(f"Processing task {task_id}")
    time.sleep(1)
    return task_id

queue = TaskQueue(max_workers=3)
for i in range(10):
    queue.submit(process_task, i)

queue.wait()

六、并发编程最佳实践

6.1 避免共享状态

# 不好的做法
shared_list = []

def append_item(item):
    shared_list.append(item)

# 好的做法
def process_item(item):
    return processed_item

6.2 使用线程安全的数据结构

import queue

thread_safe_queue = queue.Queue()

def producer():
    for i in range(10):
        thread_safe_queue.put(i)

def consumer():
    while True:
        item = thread_safe_queue.get()
        process(item)
        thread_safe_queue.task_done()

6.3 设置合理的线程/进程数

import multiprocessing

cpu_count = multiprocessing.cpu_count()

# CPU密集型任务
with ProcessPoolExecutor(max_workers=cpu_count) as executor:
    pass

# IO密集型任务
with ThreadPoolExecutor(max_workers=cpu_count * 2) as executor:
    pass

6.4 异常处理

from concurrent.futures import ThreadPoolExecutor, as_completed

def risky_operation(data):
    if data < 0:
        raise ValueError("Negative value")
    return data * 2

with ThreadPoolExecutor() as executor:
    futures = {executor.submit(risky_operation, i): i for i in [-1, 0, 1, 2]}
    
    for future in as_completed(futures):
        data = futures[future]
        try:
            result = future.result()
            print(f"{data} -> {result}")
        except ValueError as e:
            print(f"{data} raised {e}")

七、实战案例:并发文件处理系统

import os
import hashlib
from concurrent.futures import ThreadPoolExecutor

def compute_checksum(filepath):
    try:
        with open(filepath, 'rb') as f:
            hash_obj = hashlib.md5()
            while chunk := f.read(4096):
                hash_obj.update(chunk)
            return filepath, hash_obj.hexdigest()
    except Exception as e:
        return filepath, str(e)

def find_files(directory):
    files = []
    for root, dirs, filenames in os.walk(directory):
        for filename in filenames:
            files.append(os.path.join(root, filename))
    return files

def main():
    directory = "/path/to/directory"
    files = find_files(directory)
    
    print(f"Found {len(files)} files")
    
    with ThreadPoolExecutor(max_workers=os.cpu_count()) as executor:
        results = executor.map(compute_checksum, files)
        
        for filepath, checksum in results:
            print(f"{filepath}: {checksum}")

if __name__ == "__main__":
    main()

总结

并发编程是构建高性能系统的关键技术。通过本文的学习,你应该掌握了以下核心要点:

  1. 并发基础:并发与并行的区别
  2. 多线程:threading模块、线程同步、线程池
  3. 多进程:multiprocessing模块、进程间通信、进程池
  4. 性能对比:线程与进程的优缺点
  5. 实用案例:网页爬取、数据处理、任务队列
  6. 最佳实践:避免共享状态、线程安全、异常处理

作为从Python转向Rust的后端开发者,掌握并发编程对于构建高效系统至关重要。Rust的并发模型更加安全和高效,通过所有权系统避免了许多并发问题。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐