1.并发与并行

并发:单个cpu处理多个任务。各个任务交替执行一段时间

并行:多个cpu同时执行各个认为

2.多进程

进程

进程是操作系统进行资源分配的基本单位。每个进程有自己独立的一块内存空间。多进程就是在操作系统中执行多个进程。

使用multiprocessing.process创建进程

multiprocessing.Process 是 Python 中用于创建和管理进程的核心类。以下是一个完整的示例,展示如何创建并启动一个子进程。

multiprocessing.Process(group = None, target = none, name = none, args = (), kwargs = {}, *, daemon = none)
import multiprocessing
import os

def worker():
    print(f"子进程 ID: {os.getpid()}, 父进程 ID: {os.getppid()}")

if __name__ == "__main__":
    print(f"主进程 ID: {os.getpid()}")
    p = multiprocessing.Process(target=worker)
    p.start()
    p.join()
 

关键参数说明

  • target: 指定子进程要执行的函数。由run()方法来发起调用的可调用对象。
  • group:保留参数,始终设为 None,为未来扩展预留。当前版本未实现功能
  • daemon:布尔值,控制进程是否为守护进程。守护进程会随主进程退出而终止,且不允许创建子进程。默认为none则继承父进程。
  • args: 以元组形式传递参数给目标函数。
  • kwargs: 以字典形式传递关键字参数。
  • name: 为进程设置名称(可通过 p.name 查看)。

Process的属性和方法

name
进程的名称,可通过赋值修改。常用于标识进程用途

pid
返回进程的ID号,类型为整数,唯一标识一个进程。

daemon
布尔值,表示是否为守护进程。主线程退出时,守护进程会自动终止。

exitcode
进程退出时的状态码。运行中为None,正常退出为0,被信号终止为负数。

authkey
进程的身份验证密钥,用于网络连接时的安全验证,默认为os.urandom()生成。

start()
启动进程,调用后系统会创建新的进程并执行run()方法。每个进程只能调用一次。

run()
默认调用传入target的对象,定义进程执行的任务逻辑。通常通过继承Process类并重写此方法实现自定义功能。

join(timeout=None)
阻塞主进程,等待子进程结束。timeout参数指定超时时间(秒),超时后继续执行。

terminate()
强制终止子进程。可能产生资源未释放的问题,建议优先使用正常结束逻辑。

is_alive()
返回布尔值,检查进程是否仍在运行。可用于监控进程状态。

close()
释放进程资源。仅当未调用start()或进程已终止时可用。

os.getpid():获取当前进程编号

os.getppid():获取当前进程的父进程编号

示例:文件同时读写

import multiprocessing
import time


#同时对读写文件进行操作

def write_file():
    with open("text.txt","a") as f:
        while true:
            f.write("hello world\n")
            #手动将缓冲区数据刷写到文件
            f.flush()
            time.sleep(0.5)

def read_file():
    with open("text.txt","r") as f:
        while true:
            time.sleep(0.5)
            print(f.readline())


if __name__ == "__main__":
    #让当前两函数同时执行
    p1 = multiprocessing.Process(target = write_file)
    p2 = multiprocessing.Process(target = read_file)

    p1.start()
    p2.start()

若直接调用multiprocessing相关函数而不加保护条件,子进程在初始化时会再次执行模块级代码,形成循环。if __name__ == "__main__"确保代码仅在主进程中执行。

自定义Process子类创建对象

import os
import multiprocessing

class Worker(multiprocessing.Process):
    def fun(self):
        print("进程id:",os.getpid(),"\t父进程id:",os.getppid())

if __name__ == "__main__":
    #创建多个进程对象
    for i in range(5):
        p = Worker()
        p.start()

通过进程池创建进程对象

multiprocessing.Pool(
    processes=None,
    initializer=None,
    initargs=(),
    maxtasksperchild=None,
    context
)
 

processes

  • 类型:整数或 None
  • 作用:指定进程池中的工作进程数量。
  • 默认值:None,表示使用 os.cpu_count() 返回的 CPU 核心数。

initializer

  • 类型:可调用对象(如函数)或 None
  • 作用:每个工作进程启动时执行该函数。
  • 默认值:None,表示不执行任何初始化操作。

initargs

  • 类型:元组
  • 作用:传递给 initializer 函数的参数。
  • 默认值:空元组 ()

maxtasksperchild

  • 类型:整数或 None
  • 作用:每个工作进程完成指定数量的任务后会被重启,避免内存泄漏。一个工作进程在它退出或被一个新的工作进程代替之前能完成的任务数量。
  • 默认值:None,表示工作进程寿命与池对齐。

context

  • 作用:可被用于指定启动的工作进程的上下文。通常一个进程池是使用函数multiprocessing.Pool()或者上一个对象的pool()方法创建的

注意⚠️:进程池对象的方法只有创建它的进程能够调用,使用时候一般只指定process参数

进程池常用方法

apply

  • 语法:apply(func, args=(), kwds={})
  • 作用:同步调用函数,阻塞直到返回结果。

apply_async

  • 语法:apply_async(func, args=(), kwds={}, callback=None, error_callback=None)
  • 作用:异步调用函数,返回 AsyncResult 对象。

map

  • 语法:map(func, iterable, chunksize=None)
  • 作用:并行处理可迭代对象,返回结果列表。

map_async

  • 语法:map_async(func, iterable, chunksize=None, callback=None, error_callback=None)
  • 作用:异步版本的 map,返回 AsyncResult 对象。

close

  • 语法:close()
  • 作用:关闭进程池,禁止提交新任务

terminate

  • 语法:terminate()
  • 作用:立即终止所有工作进程

join

  • 语法:join()
  • 作用:阻塞主进程,等待所有工作进程结束,需在 close 或 terminate 后调用。

案例:

import os
import time
import multiprocessing

def func():
    for i in range(10)
    print(os.getpid(),i)
    time.sleep(0.5)

if __name__ == "__main__:"
    #指定进程池大小
    process_num = 5
    pool = multiprocessing.Pool(process_num)


    for p in range(process_num):
        #交任务给进程池
        """阻塞式:p.apply(func)"""
        #非阻塞式
        pool.apply_async(func)#异步
    
    pool.close()
    pool.join()
    

进程间通信

进程间不共享全局变量

子进程向传入的列表中添加元素,最终发现主进程与子进程之间的列表结果不同:

import os

def func(list1):
    for i in range(10):
        list1.append(i)
        print(os.getpid(),list1)


if __name__ == "__main__":
    list1 = []
    p1 = multiprocessing.Process(target = func, args = (list1,))
    p2 = multiprocessing.Process(target = func, args = (list1,))
    p1.start()
    p2.start()
    p1.join()
    p2.join()
    print(os.getpid(),list1)

    print(f"当前进程名{multiprocessing.current_process().name}",list1)

全局变量list1没有被p1和p2共享

使用Queue通信——可以实现多个进程之间共享

python的multiprocessing模块包装底层的机制,提供了Queue、pipes等多种方式来交换数据。

multiprocessing.Queue([maxsize])返回一个使用一个管道和少量锁和信号量实现的共享队列实例。当一个进程将一个对象放进队列中时,一个写入线程会启动并将对象从缓冲区写入管道中。默认队列是无限大小的,可以通过maxsize参数限制。

Queue常用方法

qsize():返回队列大致长度。由于多线程或者多进程上下文,这个数字是不可靠的。

empty():如果队列是空的返回True。由于多线程或者多进程上下文,这个数字是不可靠的。

full():如果队列是满的返回True

put(obj[, block[, timeout[):将obj放入队列。block=true而且timeout=none,将会阻塞当前进程,指导有空的缓冲槽。如果timeout是正数,将会阻塞了最多timeout秒之后还是没有可用的缓冲槽时,抛出queue.Full异常。反之block = false,仅当有可用缓冲槽时才放入对象,否则抛出queue.Full异常

put_nowait(obj):相当于put(obj,false)

get([block[, timeout]])):block=true而且timeout=none,将会阻塞当前进程,直到队列出现可用对象。如果timeout是正数,将会阻塞了最多timeout秒之后还是没有可用的对象时,抛出queue.empty异常。反之block = false,仅当有可用对象能够取出时才放入对象,否则抛出queue.empty异常

get_nowait(obj):相当于get(false)

案例:

import os
import multiprocessing
import random

#放数据
def func1(queue):
    while true:
        queue.put(random.randint(1,50))
        time.sleep(0.5)

#取数据
def fun2(queue):
    while true:
        print("="*queue.get())

if __name__ == "__main__":
    queue = multiprocessing.Queue()
    p1 = multiprocessing.Process(target = func1,args = (queue,))
    p2 = multiprocessing.Process(target = func2,args = (queue,))
    p1.start()
    p2.start()
    p1.join()
    p2.join()

使用进程池——要使用Manager().Queue

if __name__ == "__main__":
    qu = multiprocessing.Manager().Queue()
    pool = multiprocessing.Pool(2)
    pool.apply_async(func1,(queue,))
    pool.apply_async(func2,(queue,))
    pool.close()
    pool.join()

3.多线程

线程是处理器任务调度和执行的基本单位

一个进程至少有一个线程,也可以运行多个线程

多个线程之间可以共享数据

线程运行出错异常后,如果没有捕获,会导致整个进程崩溃

多线程是指在同一个进程中同时执行多个任务

使用threading.Thread创建线程

python提供了两个模块:_thread和threading,_thread是低级模块,threading是高级模块,对_thread的封装。绝大多数下,我们只使用threading这个高级模块

Thread的创建

threading.Thread(group = none, target = none, name = none, args = (), kwargs = {}, * , daemon = none)

group
保留参数,当前未使用,默认值为None。为未来扩展预留,通常无需指定。

target
指定线程要执行的函数或方法。必须是可调用对象(如函数名),默认None表示无操作。

name
设置线程名称,便于调试和日志记录。若不指定,Python自动生成Thread-N格式的名称(N为数字)。

args
元组形式传递给target函数的位置参数。例如args=(1, 2)会解包为target(1, 2)

kwargs
字典形式传递给target函数的关键字参数。例如kwargs={'x': 1}等价于target(x=1)

daemon
布尔值,控制线程是否为守护线程。若为True,主线程退出时会强制结束该线程;默认None继承当前线程的守护状态。

thread的属性和方法与其他常用方法

name
线程名称,用于标识线程。可通过构造函数或直接赋值设置。

ident
线程标识符。线程启动后由系统分配,未启动时为 None

daemon
布尔值,表示是否为守护线程。主线程退出时,守护线程会自动终止。需在 start() 前设置。

is_alive()
返回线程是否正在运行。线程启动后且未终止时为 True

start()
启动线程,调用 run() 方法。每个线程只能调用一次。

run()
定义线程执行逻辑。默认调用传递给构造函数的 target函数。子类可重写此方法。

join(timeout=None)
阻塞当前线程,直到目标线程结束或超时。timeout 为可选参数(秒)。

threading.current_thread()
返回当前线程对象。

threading.enumerate()
返回所有存活线程的列表,包括主线程和守护线程。

threading.active_count()
返回当前存活的线程数量。

案例:两线程分别交替打印

import time
import threading

#交替打印0000和1111
def func():
    flag = 0
    while true
        print(threading.current_thread().name, f"{flag}"*4)
        flag = flag^1 #替换1和0
        time.sleep(0.5)

if __name__ == "__main__":
    t1 = threading.Thread(Target = func, name = "线程1")
    t2 = threading.Thread(Target = func, name = "线程2")
    t1.start()
    t2.start()

 自定义Thread线程子类创建线程

import threading

class Worker(threading.Thread):
    def run(self):
        flag = 0
        while True:
            print(f"线程{threading.current_thread().name}",f"{flag}"*5)
            flag = flag^1
            time.sleep(0.6)

if __name__ == "__main__":
    t1 = Worker(name = "thread1")
    t2 = Worker(name = "thread2")
    t1.start()
    t2.start()

线程池

ThreadPoolExecutor 是concurrent.futures模块中的线程池实现,它允许我们轻松地提交任务到线程池,并管理任务的执行和结果。

线程池的创建

concurrent.futures.ThreadPoolExecutor(max_workers = none, thread_name_prefix = "",initializer = none, initargs = ())

max_workers
指定线程池中最大工作线程数。默认值为 None,此时会根据系统CPU核心数自动设置(通常为 os.cpu_count() * 5)。若显式设置为整数,则严格限制线程数量。

thread_name_prefix
为线程池中的线程设置名称前缀。默认空字符串,线程名称为 ThreadPoolExecutor-{数字}。自定义前缀可帮助调试时区分不同线程池的线程。

initializer
接收一个可调用对象(如函数),在每个线程启动时执行初始化。默认 None 表示不进行初始化。适用于线程本地资源(如数据库连接)的预配置。

initargs
以元组形式传递 initializer 函数的参数。默认 None表示无参数。需与 initializer 配合使用,例如初始化时传递配置路径。

线程池常用方法

  • submit(Runnable task):提交任务并返回Future<?>,可通过Future.get()等待任务完成。
  • shutdown():平缓关闭线程池,等待已提交任务完成,不再接受新任务。
  • shutdownNow():立即尝试终止所有任务,返回未执行的任务列表。
def func(tname):
    global word
    for i, char in enumerate(word):
        word[i] = chr(ord(char)^1)
        print(f"{tname} : {word}\n", end = "")
    return word

"""
for i, char in enumerate(word):
遍历 word 的每个字符,i 为索引,char 为当前字符。

word[i] = chr(ord(char)^1)
对每个字符执行以下操作:

ord(char):获取字符的 Unicode 码点。
^1:对码点进行按位异或(XOR)运算,参数 1 表示二进制最后一位取反。
chr(...):将运算后的码点转回字符。
将结果写回 word[i],直接修改原字符串的字符。

print(f"{tname} : {word}\n", end="")
格式化输出:

tname:函数参数,用于标识输出来源。
word:每次修改后的完整字符串。
end="":禁止换行,确保多次调用时输出紧凑。
"""
if __name__ == "__main__":
    word = list("idmmn!vnsme")
    #使用with语句来确保线程被迅速处理
    with concurrent.futures.ThreadPoolExecutor(max_worker = 3) as executor:
        future1 = executor.submit(fun,"thread1")
        future2 = executor.submit(fun,"thread2")
        future3 = executor.submit(fun,"thread3")
        word = future1.result()
        word = future2.result()
        word = future3.result()
    print("".join(word))

代码执行逻辑

  1. 初始化字符串:将字符串"idmmn!vnsme"转换为列表word,便于修改单个字符。
  2. 线程池创建:使用ThreadPoolExecutor创建最多3个线程的线程池。
  3. 线程任务提交:每个线程执行func函数,传入线程名称(如"thread1")。
  4. 字符处理:每个线程遍历word,对每个字符执行异或运算并更新列表。
  5. 结果合并:主线程等待所有线程完成,最终输出处理后的字符串。

互斥锁

线程安全问题

线程之间共享数据会存在线程安全的问题

举例:多个线程同时递增一个共享计数器,由于操作非原子性,可能导致计数错误。

import threading

counter = 0

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

threads = []
for _ in range(10):
    t = threading.Thread(target=increment)
    threads.append(t)
    t.start()
"""
创建10个线程,每个线程执行increment函数。
将线程对象存入列表threads,并立即启动线程。
"""

for t in threads:
    t.join()

print(f"Expected: 1000000, Actual: {counter}")
 

互斥锁的概念

互斥锁(Mutex,全称 Mutual Exclusion)是一种同步机制,用于控制多个线程或进程对共享资源的访问。其主要目的是防止并发访问导致的数据竞争或不一致问题。某个线程要更改共享数据时,先将其锁定,此时其他线程不能更改。知道该线程释放资源,将资源的状态变成“非锁定”,其他的线程才可以再次锁定该资源。

互斥锁的使用

可以通过threading.Lock()创建互斥锁

使用lock.acquire([blocking = True][,timeout = 1])来获取锁(blocking为true,线程会阻塞知道获取到锁;如果为false,线程立即返回。获取锁成功返回True,否则返回false。timeout为等待时间,超时未获取锁,则返回false)

使用lock.release()释放锁

def func():
    global g_num
    for _ in range(10):
        lock.acquire()
        g_num+=1
        time.sleep(0.2)
        print(f"当前线程{threading.current_thread().name}", g_num)
        lock.release()


if __name__ == "__main__":
    g_num = 0
    lock = threading.Lock()
    threads = [threading.Thread(target = func, name = "线程" + str(i)) for i in range (1,4)]
    [t.start() for t in threads]
    [t.join() for t in threads]

    print(f"当前线程{threading.current_thread().name}", g_num)
"""
g_num = 0:初始化全局变量g_num为0。
lock = threading.Lock():创建线程锁对象,用于同步线程。
thread = [...]:生成3个线程,每个线程命名为“线程1”、“线程2”、“线程3”,目标函数为func。
[t.start() for t in threads]:启动所有线程(注意变量名应为thread而非threads,此处为代码错误)。
[t.join() for t in threads]:等待所有线程执行完毕(变量名错误同上)。
print(...):主线程打印最终的g_num值。
"""
        

买票系统举例

import threading
import time

def sale_ticket():
    global ticket_num

    while True:
        lock.acquire()
        if ticket_num <= 0
            lock.release()
            break
        time.sleep(0.01)
        ticket_num -= 1
        printf(f"当前{thread.current_thread().name}卖了一张票,还剩{ticket_num}")
        lock.release()

if __name__ == "__main__":
    ticket_num = 100
    lock = threading.Lock()
    threads = [threading.Thread(Target = sale_ticket, name = "windows" + str(i)) for i in range(1,4)]
    [t.start() for t in threads]
    [t.join() for t in threads]

GIL

GIL(Global Interpreter Lock)是 Python 解释器中的一个全局锁机制,用于同步线程的执行。它确保同一时刻只有一个线程可以执行 Python 字节码这意味着在任何时间点都只能有一个线程处于执行状态,从而简化了 CPython 解释器的内存管理设计。GIL 的主要目的是保护 Python 对象免受多线程并发访问导致的内存管理问题。CPython 使用引用计数进行垃圾回收,而 GIL 避免了多个线程同时修改引用计数引发的竞争条件。

4.进程和线程对比

区别

资源分配:进程拥有独立的内存空间和系统资源,每个进程都有自己的代码段、数据段和堆栈等。而线程共享所属进程的内存空间和资源,同一进程内的线程之间可以直接访问共享内存。

开销:创建进程需要分配独立的内存、打开文件等系统资源,开销较大;创建线程只需所属进程的内存空间进行少量资源分配,开销较小;

并发行:多核心cpu环境下,进程和线程都可以异步执行但进程之间的异步是真正的异步,而线程之间的异步在单核心cpu上是通过时间片轮转实现的“伪异步”

可靠性:进程间相互隔离,一个进程的错误不会扩散,系统稳定性更高。线程共享内存,一个线程的错误可能影响其他线程,甚至导致整个进程崩溃。

通信与同步:进程间通信(IPC)需通过管道、消息队列、共享内存等机制,实现复杂且速度较慢;线程间可直接读写共享数据,但需通过锁、信号量等机制避免竞态条件,编程复杂度较高。

适合使用多线程的情况:

I/O密集型任务

程序需要频繁等待磁盘、网络或数据库响应时,多线程可避免阻塞主线程。例如Web服务器处理并发请求,或爬虫同时下载多个网页,线程在等待I/O时可切换执行其他任务,提高资源利用率。

共享数据
进程间通信(IPC)成本高于线程间通信。若任务需要频繁共享或修改同一数据,多进程可能导致性能下降。

适合使用多进程的勤情况:

CPU密集型任务
当任务需要大量计算资源且受限于单个CPU核心时,多进程能有效利用多核CPU的并行计算能力。例如科学计算、图像处理、机器学习模型训练等场景,通过多进程可显著缩短计算时间。

避免GIL限制(Python特有)
在Python中,由于全局解释器锁(GIL)的存在,多线程无法真正并行执行CPU密集型任务。多进程可以绕过GIL限制,每个进程拥有独立的Python解释器和内存空间。

任务独立性高
多进程适合处理相互独立、无需频繁共享数据的任务。例如批量文件处理、网页爬虫等场景,每个进程可独立处理部分任务,最后合并结果。

稳定性要求高
多进程中单个进程崩溃不会影响其他进程的运行,适合需要高稳定性的场景。例如长时间运行的服务,可通过多进程隔离风险。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐