python

关注公众号 jb51net

关闭
首页 > 脚本专栏 > python > Python多任务处理方法

Python多任务处理方法有哪些?多线程、异步编程怎么选?

作者:图灵追慕者

还在纠结Python并发编程怎么选?本文详解多线程、多进程、异步编程等方法,帮你根据I/O密集或CPU密集型任务轻松选择,附线程池、协程用法,避开GIL和线程安全等坑

在Python中,实现多任务处理(并发或并行处理)有多种方法。根据任务的性质(I/O密集型还是CPU密集型)以及具体需求,可以选择不同的多任务处理方法。

以下是Python中常见的多任务处理方法及其使用手法:

1.多线程(threading模块)

适用场景:I/O密集型任务,如文件读写、网络请求等。由于Python的全局解释器锁(GIL),多线程在CPU密集型任务中效果有限。

基本用法

import threading
import time
def worker(name, delay):
    print(f"线程 {name} 开始")
    time.sleep(delay)
    print(f"线程 {name} 结束")
# 创建多个线程
thread1 = threading.Thread(target=worker, args=("A", 2))
thread2 = threading.Thread(target=worker, args=("B", 3))
# 启动线程
thread1.start()
thread2.start()
# 等待所有线程完成
thread1.join()
thread2.join()
print("所有线程已完成")

常见使用手法

from concurrent.futures import ThreadPoolExecutor
import time
def worker(name, delay):
    print(f"线程 {name} 开始")
    time.sleep(delay)
    print(f"线程 {name} 结束")
    return f"结果 {name}"
with ThreadPoolExecutor(max_workers=2) as executor:
    future1 = executor.submit(worker, "A", 2)
    future2 = executor.submit(worker, "B", 3)
    print(future1.result())
    print(future2.result())
print("所有线程已完成")

2.多进程(multiprocessing模块)

适用场景:CPU密集型任务,如数学计算、数据处理等。多进程能够绕过GIL,实现真正的并行计算。

基本用法

from multiprocessing import Process
import time
def worker(name, delay):
    print(f"进程 {name} 开始")
    time.sleep(delay)
    print(f"进程 {name} 结束")
if __name__ == "__main__":
    process1 = Process(target=worker, args=("A", 2))
    process2 = Process(target=worker, args=("B", 3))
    process1.start()
    process2.start()
    process1.join()
    process2.join()
    print("所有进程已完成")

常见使用手法

from concurrent.futures import ProcessPoolExecutor
import time
def worker(name, delay):
    print(f"进程 {name} 开始")
    time.sleep(delay)
    print(f"进程 {name} 结束")
    return f"结果 {name}"
if __name__ == "__main__":
    with ProcessPoolExecutor(max_workers=2) as executor:
        future1 = executor.submit(worker, "A", 2)
        future2 = executor.submit(worker, "B", 3)
        print(future1.result())
        print(future2.result())
    print("所有进程已完成")

3.异步编程(asyncio模块)

适用场景:高并发I/O密集型任务,如网络服务器、爬虫等。适用于需要处理大量并发连接但每个连接处理时间较短的场景。

基本用法

import asyncio
async def worker(name, delay):
    print(f"任务 {name} 开始")
    await asyncio.sleep(delay)
    print(f"任务 {name} 结束")
    return f"结果 {name}"
async def main():
    tasks = [
        asyncio.create_task(worker("A", 2)),
        asyncio.create_task(worker("B", 3))
    ]
    results = await asyncio.gather(*tasks)
    for result in results:
        print(result)
# 运行异步主函数
asyncio.run(main())

常见使用手法

import asyncio
import aiohttp
async def fetch(session, url):
    async with session.get(url) as response:
        return await response.text()
async def main():
    urls = [
        "http://example.com",
        "http://python.org",
        "http://github.com"
    ]
    async with aiohttp.ClientSession() as session:
        tasks = [fetch(session, url) for url in urls]
        pages_content = await asyncio.gather(*tasks)
        for content in pages_content:
            print(f"下载了 {len(content)} 字节")
asyncio.run(main())

4.并发执行(concurrent.futures模块)

适用场景:简化多线程和多进程的使用,适用于需要高层次并发管理的场景。

基本用法

from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def worker(name, delay):
    print(f"线程 {name} 开始")
    time.sleep(delay)
    print(f"线程 {name} 结束")
    return f"结果 {name}"
with ThreadPoolExecutor(max_workers=3) as executor:
    futures = [executor.submit(worker, f"任务-{i}", i) for i in range(1, 4)]
    for future in as_completed(futures):
        print(future.result())
print("所有线程已完成")
from concurrent.futures import ProcessPoolExecutor, as_completed
import time
def compute(n):
    print(f"计算 {n} 开始")
    time.sleep(2)
    result = n * n
    print(f"计算 {n} 结束")
    return result
with ProcessPoolExecutor(max_workers=2) as executor:
    futures = [executor.submit(compute, i) for i in range(1, 5)]
    for future in as_completed(futures):
        print(f"结果: {future.result()}")
print("所有计算已完成")

5.协程库(如gevent、eventlet)

适用场景:需要轻量级并发,同时利用协程进行I/O操作的场景。

基本用法(以gevent为例)

import gevent
from gevent import monkey
import time
# 打补丁,使标准库中的阻塞操作变为非阻塞
monkey.patch_all()
def worker(name, delay):
    print(f"协程 {name} 开始")
    gevent.sleep(delay)
    print(f"协程 {name} 结束")
def main():
    tasks = [
        gevent.spawn(worker, "A", 2),
        gevent.spawn(worker, "B", 3),
        gevent.spawn(worker, "C", 1)
    ]
    gevent.joinall(tasks)
    print("所有协程已完成")
if __name__ == "__main__":
    main()

常见注意事项

1.线程安全:在多线程环境下,共享数据需要使用锁(Lock)或其他同步机制,避免数据竞争和不一致。

import threading
lock = threading.Lock()
shared_resource = 0
def increment():
    global shared_resource
    with lock:
        shared_resource += 1
        print(f"共享资源: {shared_resource}")
threads = [threading.Thread(target=increment) for _ in range(10)]
for t in threads:
    t.start()
for t in threads:
    t.join()

2.进程间通信(IPC):使用QueuePipe等机制在多进程间传递数据。

from multiprocessing import Process, Queue
def worker(q, n):
    q.put(n * n)
if __name__ == "__main__":
    q = Queue()
    processes = [Process(target=worker, args=(q, i)) for i in range(5)]
    for p in processes:
        p.start()
    for p in processes:
        p.join()
    results = []
    while not q.empty():
        results.append(q.get())
    print("结果:", results)

3.异常处理:在多任务处理中,确保捕获和处理可能的异常,避免任务崩溃影响整体程序。

from concurrent.futures import ThreadPoolExecutor
def worker(n):
    if n == 3:
        raise ValueError("错误发生")
    return n * n
with ThreadPoolExecutor(max_workers=2) as executor:
    futures = [executor.submit(worker, i) for i in range(5)]
    for future in futures:
        try:
            result = future.result()
            print(f"结果: {result}")
        except Exception as e:
            print(f"任务出错: {e}")

4.资源管理:确保正确关闭线程池、进程池等,避免资源泄漏。使用上下文管理器(with语句)可以简化资源管理。

5.选择合适的方法:根据任务的性质选择合适的多任务处理方法。例如,I/O密集型任务使用多线程或异步编程,CPU密集型任务使用多进程。

总结

Python提供了多种多任务处理的方法,每种方法都有其适用的场景和特点。

选择合适的方法不仅能提高程序的执行效率,还能简化代码复杂度。

理解不同方法的工作原理和使用技巧,对于编写高效且稳定的并发程序至关重要。

以上为个人经验,希望能给大家一个参考,也希望大家多多支持脚本之家。

您可能感兴趣的文章:
阅读全文