Python多线程与多进程:并发编程实战
·
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()
总结
并发编程是构建高性能系统的关键技术。通过本文的学习,你应该掌握了以下核心要点:
- 并发基础:并发与并行的区别
- 多线程:threading模块、线程同步、线程池
- 多进程:multiprocessing模块、进程间通信、进程池
- 性能对比:线程与进程的优缺点
- 实用案例:网页爬取、数据处理、任务队列
- 最佳实践:避免共享状态、线程安全、异常处理
作为从Python转向Rust的后端开发者,掌握并发编程对于构建高效系统至关重要。Rust的并发模型更加安全和高效,通过所有权系统避免了许多并发问题。
更多推荐



所有评论(0)