一、多进程编程核心概念

1. 进程的本质理解

进程(Process) 是操作系统进行资源分配的基本单位,每个进程拥有:

  • 独立的内存空间
  • 独立的代码执行环境
  • 独立的系统资源(文件描述符、网络连接等)
操作系统
进程1
进程2
内存空间
代码段
数据段
堆栈
内存空间
代码段
数据段
堆栈

2. 多进程 vs 多线程

特性 多进程 多线程
内存隔离 完全隔离 共享内存空间
创建开销 高(需复制父进程) 低(共享父进程资源)
通信成本 高(需IPC机制) 低(直接共享内存)
容错性 高(进程崩溃不影响其他) 低(线程崩溃影响整个进程)
GIL影响 无(真正并行) 有(伪并行)
适用场景 CPU密集型计算 I/O密集型操作

3. Python GIL机制解析

全局解释器锁(GIL) 是CPython解释器的限制:

  • 单进程中多线程无法真正并行
  • 多进程可完全规避GIL限制
  • CPU密集型任务首选多进程
CPU核心1
进程1
进程2
CPU核心2
进程3
进程4

4. 多进程核心优势

  • 真正并行:利用多核CPU能力
  • 内存安全:进程间内存隔离
  • 资源控制:精确分配CPU资源
  • 稳定性高:单进程崩溃不影响整体

二、multiprocessing模块基础

1. 模块架构概览

import multiprocessing

# 主要组件:
# - Process:进程类
# - Queue:进程间通信队列
# - Pipe:进程间通信管道
# - Pool:进程池
# - Manager:共享数据管理
# - Value/Array:共享内存
# - Lock/RLock:进程同步

2. 创建进程的两种方法

方法1:函数式创建
import multiprocessing
import os

def worker(num):
    print(f'子进程 {num} PID={os.getpid()}')
    return num * num

if __name__ == '__main__':
    processes = []
    
    for i in range(3):
        p = multiprocessing.Process(target=worker, args=(i,))
        processes.append(p)
        p.start()
        print(f'主进程启动子进程 PID={p.pid}')
    
    for p in processes:
        p.join()
        print(f'子进程 {p.pid} 结束')
方法2:类继承式创建
import multiprocessing
import os

class MyProcess(multiprocessing.Process):
    def __init__(self, num):
        super().__init__()
        self.num = num
        
    def run(self):
        print(f'子进程 {self.num} PID={os.getpid()}')
        result = self.num * self.num
        print(f'计算结果: {result}')

if __name__ == '__main__':
    processes = []
    
    for i in range(3):
        p = MyProcess(i)
        processes.append(p)
        p.start()
    
    for p in processes:
        p.join()
        print(f'子进程 {p.pid} 结束')

三、multiprocessing常用方法详解

1. Process类核心方法

方法 描述 参数 返回值
start() 启动进程
run() 进程执行体
join(timeout=None) 等待进程结束 timeout: 超时时间 进程是否存活
terminate() 终止进程
is_alive() 检查进程状态 布尔值
name 进程名称 可读写属性 字符串
pid 进程ID 只读属性 整数
daemon 守护进程标志 可读写属性 布尔值

2. 进程间通信(IPC)

(1) Queue进程队列
import multiprocessing
import time

def producer(q):
    for i in range(5):
        item = f'产品{i}'
        q.put(item)
        print(f'生产者 放入: {item}')
        time.sleep(0.5)

def consumer(q):
    while True:
        item = q.get()
        if item == 'END':
            break
        print(f'消费者 取出: {item}')
        time.sleep(0.3)

if __name__ == '__main__':
    # 创建进程安全队列
    queue = multiprocessing.Queue(maxsize=3)
    
    # 创建进程
    p1 = multiprocessing.Process(target=producer, args=(queue,))
    p2 = multiprocessing.Process(target=consumer, args=(queue,))
    
    p1.start()
    p2.start()
    
    # 等待生产者结束
    p1.join()
    
    # 发送结束信号
    queue.put('END')
    
    # 等待消费者结束
    p2.join()
    print('所有任务完成')
(2) Pipe进程管道
import multiprocessing

def sender(conn):
    messages = ['Hello', 'World', 'Python']
    for msg in messages:
        conn.send(msg)
        print(f'发送: {msg}')
        response = conn.recv()
        print(f'收到响应: {response}')
    conn.send('END')
    conn.close()

def receiver(conn):
    while True:
        msg = conn.recv()
        if msg == 'END':
            conn.send('确认结束')
            break
        print(f'接收: {msg}')
        conn.send(f'已收到: {msg}')

if __name__ == '__main__':
    # 创建管道
    parent_conn, child_conn = multiprocessing.Pipe()
    
    # 创建进程
    p1 = multiprocessing.Process(target=sender, args=(parent_conn,))
    p2 = multiprocessing.Process(target=receiver, args=(child_conn,))
    
    p1.start()
    p2.start()
    
    p1.join()
    p2.join()
    print('双向通信完成')

3. 进程池(Pool)管理

import multiprocessing
import time

def cpu_intensive_task(n):
    """CPU密集型任务"""
    print(f'开始计算 {n}²')
    result = 0
    for i in range(n * 1000000):
        result += i % 256
    return result

if __name__ == '__main__':
    # 创建进程池
    with multiprocessing.Pool(processes=4) as pool:
        # 方法1: apply_async 异步提交
        results_async = []
        for i in range(1, 6):
            result = pool.apply_async(cpu_intensive_task, (i,))
            results_async.append(result)
        
        # 获取异步结果
        print('异步结果:')
        for res in results_async:
            print(res.get(timeout=10))  # 设置超时
        
        # 方法2: map 批量处理
        print('\nmap批量处理:')
        inputs = [10, 20, 30, 40]
        outputs = pool.map(cpu_intensive_task, inputs)
        print(outputs)
        
        # 方法3: map_async 异步批量处理
        print('\nmap_async异步批量处理:')
        result_async = pool.map_async(cpu_intensive_task, [5, 15, 25])
        print(result_async.get())  # 获取结果
    
    print('所有任务完成')

4. 共享内存与同步

(1) Value & Array
import multiprocessing
import time

def increment(counter, lock):
    """安全递增计数器"""
    for _ in range(100000):
        with lock:
            counter.value += 1

if __name__ == '__main__':
    # 创建共享计数器
    counter = multiprocessing.Value('i', 0)
    
    # 创建锁
    lock = multiprocessing.Lock()
    
    # 创建进程
    processes = []
    for _ in range(4):
        p = multiprocessing.Process(target=increment, args=(counter, lock))
        processes.append(p)
        p.start()
    
    # 等待所有进程完成
    for p in processes:
        p.join()
    
    print(f'最终计数: {counter.value} (期望: 400000)')
(2) Manager共享字典
import multiprocessing

def worker(shared_dict, key, value):
    """修改共享字典"""
    shared_dict[key] = value
    print(f'进程 {multiprocessing.current_process().name} 设置 {key}={value}')

if __name__ == '__main__':
    with multiprocessing.Manager() as manager:
        # 创建共享字典
        shared_dict = manager.dict()
        
        processes = []
        items = [('A', 1), ('B', 2), ('C', 3), ('D', 4)]
        
        # 创建并启动进程
        for key, value in items:
            p = multiprocessing.Process(
                target=worker,
                args=(shared_dict, key, value)
            )
            processes.append(p)
            p.start()
        
        # 等待所有进程完成
        for p in processes:
            p.join()
        
        # 输出最终结果
        print(f'共享字典内容: {dict(shared_dict)}')

四、Python多进程实战案例

1. 科学计算:矩阵并行乘法

import multiprocessing
import numpy as np
import time

def matrix_multiply(args):
    """子矩阵乘法任务"""
    A, B, start_row, end_row = args
    result = np.zeros((end_row - start_row, B.shape[1]))
    
    for i in range(start_row, end_row):
        for j in range(B.shape[1]):
            for k in range(A.shape[1]):
                result[i - start_row, j] += A[i, k] * B[k, j]
    
    return (start_row, end_row, result)

def parallel_matrix_multiply(A, B, num_processes=None):
    """并行矩阵乘法"""
    if A.shape[1] != B.shape[0]:
        raise ValueError("矩阵维度不匹配")
    
    if num_processes is None:
        num_processes = multiprocessing.cpu_count()
    
    # 计算每个进程处理的行数
    rows_per_process = A.shape[0] // num_processes
    tasks = []
    
    # 准备任务
    for i in range(num_processes):
        start_row = i * rows_per_process
        end_row = (i + 1) * rows_per_process if i < num_processes - 1 else A.shape[0]
        tasks.append((A, B, start_row, end_row))
    
    # 并行计算
    with multiprocessing.Pool(processes=num_processes) as pool:
        results = pool.map(matrix_multiply, tasks)
    
    # 合并结果
    C = np.zeros((A.shape[0], B.shape[1]))
    for start_row, end_row, partial in results:
        C[start_row:end_row, :] = partial
    
    return C

if __name__ == '__main__':
    # 创建大型矩阵 (1000x1000)
    A = np.random.rand(1000, 1000)
    B = np.random.rand(1000, 1000)
    
    # 串行计算
    print("串行计算开始...")
    start = time.time()
    C_serial = np.dot(A, B)
    print(f"串行计算耗时: {time.time() - start:.4f}秒")
    
    # 并行计算
    print("\n并行计算开始...")
    start = time.time()
    C_parallel = parallel_matrix_multiply(A, B, num_processes=8)
    print(f"并行计算耗时: {time.time() - start:.4f}秒")
    
    # 验证结果
    assert np.allclose(C_serial, C_parallel), "结果不一致"
    print("结果验证通过")

2. 数据处理:日志分析系统

import multiprocessing
import re
import time
from collections import defaultdict

def log_processor(chunk, pattern):
    """日志处理函数"""
    print(f"进程 {multiprocessing.current_process().name} 处理 {len(chunk)} 行日志")
    
    # 编译正则表达式
    regex = re.compile(pattern)
    stats = defaultdict(int)
    
    for line in chunk:
        # 匹配IP地址
        if re.search(r'(\d+\.\d+\.\d+\.\d+)', line):
            stats['ip_count'] += 1
        
        # 匹配HTTP状态码
        match = regex.search(line)
        if match:
            status = match.group(1)
            stats[f'status_{status}'] += 1
        
        # 错误统计
        if 'ERROR' in line:
            stats['errors'] += 1
        elif 'WARN' in line:
            stats['warnings'] += 1
    
    return dict(stats)

def analyze_logs(file_path, pattern, num_processes=4, chunk_size=5000):
    """并行日志分析"""
    # 读取日志文件
    with open(file_path, 'r') as f:
        lines = f.readlines()
    
    print(f"共读取 {len(lines)} 行日志")
    
    # 分割日志块
    chunks = [lines[i:i+chunk_size] for i in range(0, len(lines), chunk_size)]
    print(f"分割为 {len(chunks)} 个任务块")
    
    # 并行处理
    with multiprocessing.Pool(processes=num_processes) as pool:
        results = pool.starmap(log_processor, [(chunk, pattern) for chunk in chunks])
    
    # 合并结果
    total_stats = defaultdict(int)
    for result in results:
        for key, value in result.items():
            total_stats[key] += value
    
    return dict(total_stats)

if __name__ == '__main__':
    # 日志文件路径
    log_file = "server.log"
    
    # 正则模式匹配HTTP状态码
    status_pattern = r'HTTP\/\d\.\d\" (\d{3})'
    
    # 分析日志
    results = analyze_logs(log_file, status_pattern, num_processes=8)
    
    # 打印结果
    print("\n=== 日志分析结果 ===")
    print(f"总访问IP数: {results.get('ip_count', 0)}")
    print(f"错误日志数: {results.get('errors', 0)}")
    print(f"警告日志数: {results.get('warnings', 0)}")
    
    print("\nHTTP状态码分布:")
    for key, value in results.items():
        if key.startswith('status_'):
            print(f"  {key.split('_')[1]}: {value}")

3. 混合编程模型:进程+线程池

import multiprocessing
from concurrent.futures import ThreadPoolExecutor
import time
import requests

def fetch_url(url):
    """网络请求函数"""
    try:
        response = requests.get(url, timeout=5)
        return {
            'url': url,
            'status': response.status_code,
            'size': len(response.text),
            'success': True
        }
    except Exception as e:
        return {
            'url': url,
            'error': str(e),
            'success': False
        }

def site_checker(urls, results):
    """站点检查进程"""
    print(f"检查进程 PID={multiprocessing.current_process().pid} 启动")
    
    # 在进程内创建线程池
    with ThreadPoolExecutor(max_workers=10) as executor:
        # 提交所有任务
        futures = {executor.submit(fetch_url, url): url for url in urls}
        
        # 处理结果
        for future in futures:
            result = future.result()
            if result['success']:
                print(f"成功: {result['url']} 状态码: {result['status']}")
            else:
                print(f"失败: {result['url']} 错误: {result['error']}")
            
            # 存储结果
            results.append(result)

if __name__ == '__main__':
    # 目标URL列表(100个网站)
    top_sites = [...]  # 替换为实际URL列表
    
    # 创建共享结果列表
    with multiprocessing.Manager() as manager:
        results = manager.list()
        
        # 分割任务到4个进程
        num_processes = 4
        urls_per_process = len(top_sites) // num_processes
        processes = []
        
        # 启动进程
        for i in range(num_processes):
            start_idx = i * urls_per_process
            end_idx = start_idx + urls_per_process if i < num_processes - 1 else len(top_sites)
            process_urls = top_sites[start_idx:end_idx]
            
            p = multiprocessing.Process(
                target=site_checker,
                args=(process_urls, results)
            )
            processes.append(p)
            p.start()
        
        # 等待所有进程完成
        for p in processes:
            p.join()
        
        # 分析结果
        success = sum(1 for r in results if r['success'])
        failure = len(results) - success
        
        print(f"\n检查完成: 成功={success} 失败={failure}")

五、最佳实践与性能优化

1. 多进程编程黄金法则

  1. 资源隔离原则

    • 避免共享状态,优先使用IPC
    • 必须共享时使用Manager或共享内存
    • 谨慎处理文件描述符
  2. 进程生命周期管理

    # 安全进程管理模板
    processes = []
    try:
        for task in tasks:
            p = Process(target=worker, args=(task,))
            p.start()
            processes.append(p)
        
        for p in processes:
            p.join(timeout=60)  # 设置超时
    finally:
        # 异常时终止所有进程
        for p in processes:
            if p.is_alive():
                p.terminate()
    
  3. 性能优化策略

    • 避免小任务(进程启动开销>任务开销)
    • 批量处理数据减少IPC次数
    • CPU密集型任务:进程数=CPU核心数
    • I/O密集型任务:增加进程内线程数

2. 跨平台兼容性方案

import multiprocessing
import platform
import os

def safe_process_creation():
    """跨平台安全创建进程"""
    system = platform.system()
    
    # 设置合适的启动方法
    if system == 'Windows':
        # Windows只支持spawn
        ctx = multiprocessing.get_context('spawn')
    elif system == 'Darwin':  # macOS
        # macOS避免fork问题
        ctx = multiprocessing.get_context('forkserver')
    else:  # Linux/Unix
        ctx = multiprocessing.get_context('fork')
    
    # 创建进程
    p = ctx.Process(target=worker_function)
    p.start()
    p.join()

# 使用
if __name__ == '__main__':
    safe_process_creation()

3. 常见问题解决方案

问题 现象 解决方案
僵尸进程 进程结束但未回收 使用join()或terminate()
死锁 进程相互等待 避免嵌套锁,设置超时
内存泄漏 内存持续增长 使用进程池复用资源
序列化错误 无法pickle对象 使用Manager代理复杂对象
文件描述符泄漏 打开文件未关闭 使用with语句管理资源

通过掌握这些核心知识,您将能够高效利用Python多进程能力,突破GIL限制,构建高性能并行应用系统。

Logo

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

更多推荐