
第 13 章 进程与线程13.1并发与并行并发单个 CPU 处理多个任务。各个任务交替执行一段时间。并行多个 CPU 同时执行多个任务。13.2多进程13.2.1什么是进程进程是操作系统进行资源分配的基本单位。操作系统中一个正在运行的程序或软件就是一个进程。每个进程都有自己独立的一块内存空间。一个进程崩溃后在保护模式下不会对其他进程产生影响。多进程是指在操作系统中同时运行多个程序。13.2.2 使用multiprocessing.Process创建进程是非守护进程Unix/Linux操作系统提供了一个** os.fork()系统调用它非常特殊。普通的函数调用调用一次返回一次但是fork() **调用一次返回两次因为操作系统自动把当前进程父进程复制了一份子进程然后分别在父进程和子进程内返回。Windows 中没有** fork() **调用不过Python提供了一个跨平台的多进程模块multiprocessing。**multiprocessing模块提供了一个Process **类来代表一个进程对象。1Process 的创建multiprocessing.Process(groupNone, targetNone, nameNone, args(), kwargs{}, *, daemonNone)group应当始终为None它的存在仅是为了与threading.Thread兼容。target由run()方法来发起调用的可调用对象默认为None。**name**进程名称默认为None则自动分配。**args**针对目标调用的参数元组。**kwargs**针对目标调用的关键字参数字典。daemon是否为守护进程True或False。默认为None则继承父进程。2Process 的属性和方法与其他常用方法**name**获取进程名称。**pid**获取进程号。**daemon**判断或设置进程是否为守护进程。**exitcode**获取子进程的退出状态码。**start()**启动进程调用传入target的对象。start()只能被调用一次。run()******默认调用传入 target的对象如果子类化了 Process可以重写此方法来自定义行为**。**join([timeout])**阻塞主进程直到子进程结束或超时。timeout参数可选意为阻塞多少秒。**terminate()**强制终止子进程。**kill()**杀死进程与terminate()类似但更彻底。**is_alive()**检查进程是否仍在运行。**os.getpid()**获取当前进程编号。**os.getppid()**获取当前进程的父进程编号。案例同时读写文件import multiprocessing import time def write_file(): print(__name__,1~~~~~~~~~~) with open(output.txt, w,encodingutf-8) as f: while True: f.write(hello world\n) # 将缓冲区数据刷写到文件中 f.flush() time.sleep(0.1) def read_file(): print(__name__,2~~~~~~~~~~) with open(output.txt,r,encodingutf-8) as f: while True: time.sleep(0.1) #添加一点时间,先写再读否则如果还没有该文件容易读不到 print(f.readline()) # 在 windows 中通过 multiprocessing.Process创建进程,__name__ __main__ 必须要加 if __name__ __main__: # 创建进程 p1 multiprocessing.Process(targetwrite_file) #任务本质就是开启了一个新的解释器 p2 multiprocessing.Process(targetread_file) #启动进程 p1.start() p2.start()注意点1.创建好的进程需要启动start()2.multiprocessing.Process 创建的进程 本质就是开了一个新的解释器模块名称会改为__**mp_main__****,**然后再执行一遍代码所以如果不添加判断会循环创建进程因此在windows下multiprocessing.Process创建进程,__name__ “__main__” 必须要加3.在执行写的进程的时候任务不会立刻写到txt上而是先存到缓冲区中为了防止后续读的时候读取不到数据需要刷写到文件中f.flush()4.读取文件的时候如果还没有创建该文件后续每次写的内容都在光标之前就会读取不到因此添加一个sleep睡眠暂缓读取5.实际项目中不采用sleep方式只用于测试6.无法确认哪个子进程先进行13.2.3 自定义Process 子类创建进程是非守护进程import os import multiprocessing class Worker(multiprocessing.Process): def run(self): print(进程id, os.getpid(), \t父进程id, os.getppid()) if __name__ __main__: for i in range(5): p Worker(name进程 str(i)) p.start()multiprocessing.Process 中的multiprocessing 是包context是模块而Process只是个类为什么可以包.类的形式引用呢正常应是包下面是模块文件应该是multiprocessing.模块名.类__all__ [x for x in dir(context._default_context) if not x.startswith(_)] globals().update((name, getattr(context._default_context, name)) for name in __all__)原因init.py做了导入导出multiprocessing的__init__模块中导入了所有的模块内容并进行了处理13.2.4 进程池是守护进程当需要启动大量子进程时可以使用进程池。1)进程池的创建multiprocessing.Pool([processes[,initializer[,initargs[,maxtasksperchild[,context]]]]])processes要使用的工作进程数量。如果processes为None则使用os.cpu_count()所返回的数值。**initializer**如果不为None则每个工作进程将会在启动时调用initializer(*initargs)。**maxtasksperchild**一个工作进程在它退出或被一个新的工作进程代替之前能完成的任务数量为了释放未使用的资源。默认的maxtasksperchild是None意味着工作进程寿与池齐。**context**可被用于指定启动的工作进程的上下文。通常一个进程池是使用函数multiprocessing.Pool()或者一个上下文对象的Pool()方法创建的。注意进程池对象的方法只有创建它的进程能够调用。使用时一般只指定processes参数。2)进程池的常用方法**apply(func[, args[, kwds]])**使用args参数以及kwds命名参数同步调用func, 在返回结果前阻塞。另外func只会在一个进程池中的一个工作进程中执行。**apply_async(func[, args[, kwds[, callback[, error_callback]]]])**使用args参数以及kwds命名参数异步调用func并立即返回一个AsyncResult对象不会阻塞。可以通过callback获取结果和通过error_callback处理异常。**close()**阻止后续任务提交到进程池当所有任务执行完成后工作进程会退出。**terminate()**不必等待未完成的任务立即停止工作进程。当进程池对象被垃圾回收时会立即调用terminate()。**join()**阻塞主进程等待工作进程结束。调用join()前必须先调用close()或者terminate()。**注意**从3.8版本后进程池的进程默认是守护进程所以需 join() 确保主进程等待。3)案例import os import time import multiprocessing # 打印10个数字,每次间隔0.5秒 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): # 阻塞式 # pool.apply(func) # 非阻塞式 pool.apply_async(func) # 异步 pool.close() # 堵塞后续任务提交到进程池, pool.join() #堵塞主线程 print(end)13.2.5 进程间通信1)进程间不共享全局变量子进程向传入的列表中添加元素最终发现主进程与子进程之间的列表结果不同import os import multiprocessing # 向list1中添加10个元素 def func(list1): for i in range(10): list1.append(i) print(os.getpid(), list1) if __name__ __main__: list1 [] p1 multiprocessing.Process(targetfunc, args(list1,)) p2 multiprocessing.Process(targetfunc, args(list1,)) p1.start() p2.start() p1.join() p2.join() print(os.getpid(), list1)args参数需要是元组的形式2)使用 Queue 通信Python的multiprocessing模块包装了底层的机制提供了Queue、Pipes等多种方式来交换数据。multiprocessing.Queue([maxsize])返回一个使用一个管道和少量锁和信号量实现的共享队列先进先出实例。当一个进程将一个对象放进队列中时一个写入线程会启动并将对象从缓冲区写入管道中。默认队列是无限大小的可以通过maxsize参数限制。1)Queue的常用方法qsize()返回队列的大致长度。由于多线程或者多进程的上下文这个数字是不可靠的。**empty()**如果队列是空的返回True。由于多线程或多进程的环境该状态是不可靠的。**full()**如果队列是满的返回True。由于多线程或多进程的环境该状态是不可靠的。**put(obj[, block[, timeout]])******将obj放入队列。如果可选参数block是True默认值而且timeout是None默认值将会阻塞当前进程直到有空的缓冲槽。如果timeout是正数将会在阻塞了最多timeout秒之后还是没有可用的缓冲槽时抛出queue.Full异常。反之block是False时仅当有可用缓冲槽时才放入对象否则抛出queue.Full异常在这种情形下timeout参数会被忽略。**put_nowait(obj)**相当于put(obj, False)。**get([block[, timeout]])******从队列中取出并返回对象。如果可选参数block是True默认值而且timeout是None默认值将会阻塞当前进程直到队列中出现可用的对象。如果timeout是正数将会在阻塞了最多timeout秒之后还是没有可用的对象时抛出queue.Empty异常。反之block是False时仅当有可用对象能够取出时返回否则抛出queue.Empty异常在这种情形下timeout参数会被忽略。**get_nowait()**相当于get(False)。3)案例:两个进程分别读写Queue(数据共享)方式1:import time import random import multiprocessing import os *# 间隔随机时间向queue中放入随机数* def func1(queue): while True: rand_num random.randint(1,50) queue.put(rand_num) print(f进程{os.getpid()}向队列中放入了元素{rand_num}) time.sleep(random.random()) *# 从queue中取出数据* def func2(queue): while True: num queue.get() print(f进程{os.getpid()}从队列中取出了元素{num}) if __name__ __main__: queue multiprocessing.Queue() p1 multiprocessing.Process(targetfunc1, args(queue,)) p2 multiprocessing.Process(targetfunc2, args(queue,)) p1.start() p2.start()使用Manager().Queue()import time import random import multiprocessing import os *# 间隔随机时间向queue中放入随机数* def func1(queue): while True: rand_num random.randint(1,50) queue.put(rand_num) print(f进程{os.getpid()}向队列中放入了元素{rand_num}) time.sleep(random.random()) *# 从queue中取出数据* def func2(queue): while True: num queue.get() print(f进程{os.getpid()}从队列中取出了元素{num}) if __name__ __main__: queue multiprocessing.Manager().Queue() p1 multiprocessing.Process(targetfunc1, args(queue,)) p2 multiprocessing.Process(targetfunc2, args(queue,)) p1.start() p2.start() p1.join() p2.join()方式2:进程池之间使用 Manager().Queue 通信import time import random import multiprocessing import os *# 间隔随机时间向queue中放入随机数* def func1(queue): while True: rand_num random.randint(1,50) queue.put(rand_num) print(f进程{os.getpid()}向队列中放入了元素{rand_num}) time.sleep(random.random()) *# 从queue中取出数据* def func2(queue): while True: num queue.get() print(f进程{os.getpid()}从队列中取出了元素{num}) if __name__ __main__: queue multiprocessing.Manager().Queue() pool multiprocessing.Pool(2) pool.apply_async(func1, args(queue,)) pool.apply_async(func2, args(queue,)) pool.close() pool.join()注意multiprocessing.Queue存在兼容性问题如果要使用进程池可以使用Mananger().Queue13.3 多线程线程是处理器任务调度和执行的基本单位。一个进程至少有一个线程也可以运行多个线程。多个线程之间可共享数据。线程运行出错异常后如果没有捕获会导致整个进程崩溃。多线程是指在同一进程中同时执行多个任务。13.3.1使用threading.Thread创建线程Python的标准库提供了两个模块_thread和threading_thread是低级模块threading是高级模块对_thread进行了封装。绝大多数情况下我们只需要使用threading这个高级模块。1)Thread 的创建threading.Thread(groupNone, targetNone, nameNone, args(), kwargs{}, *, daemonNone)**group**应为None保留给将来实现ThreadGroup类的扩展使用。**target**用于run()方法调用的可调用对象。默认是None表示不需要调用任何方法。**name**线程名称。 在默认情况下会以 “Thread-N” 的形式构造唯一名称其中 N 为一个较小的十进制数值或是 “Thread-N (target)” 的形式其中 “target” 为target.name如果指定了target参数的话。**args**用于发起调用目标函数的参数列表或元组。 默认为 ()。**kwargs**用于调用目标函数的关键字参数字典。默认是 {}。daemonTrue或False来设置该线程是否为守护模式。如果是None默认值线程将继承当前线程的守护模式属性。2)Thread 的属性和方法与其他常用方法name线程的名称。daemon线程是否为守护线程。ident线程标识符。**native_id**此线程的线程idtid由 OS内核分配。start()启动线程调用线程的 run() 方法。run()定义线程的行为默认调用传入的 target 对象。join([timeoutNone])阻塞主线程直到当前线程运行完成或达到超时时间。is_alive()线程是否在运行。threading.enumerate()查看都有哪些线程。threading.current_thread()返回当前线程实例。3)两线程分别交替打印import threading import time # 交替打印 00000 和 11111 def func(): flag 0 while True: print(threading.current_thread().name, f{flag}*5) flag flag ^ 1 #替换0 和1,异或操作 time.sleep(0.5) if __name__ __main__: t1 threading.Thread(targetfunc,namet1) t2 threading.Thread(targetfunc,namet2) t1.start() t2.start() print(~~~~~主线程~~~~~~~~)13.3.2 自定义Thread子类创建线程import threading import time class Worker123(threading.Thread): def run(self): flag 0 while True: print(threading.current_thread().name, f{flag}*5) flag flag ^ 1 time.sleep(0.5) if __name__ __main__: t1 Worker123(name线程1) t2 Worker123(name线程2) t1.start() t2.start() print(~~~~~主线程~~~~~~~~)13.3.3线程池ThreadPoolExecutor是concurrent.futures模块中的线程池实现它允许我们轻松地提交任务到线程池并管理任务的执行和结果。1)线程池的创建concurrent.futures.ThreadPoolExecutor(max_workersNone, thread_name_prefix, initializerNone, initargs())**max_workers**线程池的最大线程数默认取决于系统资源。**thread_name_prefix**线程名称前缀。**initializer**可选的初始化函数。**initargs**传递给初始化函数的参数。2)线程池的常用方法**submit(fn, *args, **kwargs)****提交一个任务到线程池返回一个Future对象。可使用Future.result() **获取任务结果。map(func, *iterables, timeoutNone, chunksize1)类似于内置的map()函数但在线程池中并行执行。Iterables为可迭代对象传递给目标函数**。chunksize **对 **ThreadPoolExecutor **没有效果。**shutdown(waitTrue, cancel_futuresFalse)**关闭线程池等待所有任务完成。wait表示是否等待线程池中的所有线程完成任务。**cancel_futures **表示是否取消尚未开始的任务。3)案例3个线程每个线程都将字符列表中的每个字符与 1 异或。import concurrent.futures 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 if __name__ __main__: word list(idmmn!vnsme) *# 使用 with 语句来确保线程被迅速清理* * *with concurrent.futures.ThreadPoolExecutor(max_workers3) as executor: future1 executor.submit(func, 线程1) *# 如果不输出结果可以不接收* * *future2 executor.submit(func, 线程2) future3 executor.submit(func, 线程3) *# word future1.result()* * # word future2.result()* * # word future3.result()* print(.join(word)) *# hello world*13.4 线程安全问题比如下面这段代码3个线程每个线程都将g_num 1 十次import time import threading def func(): global g_num for _ in range(10): tmp g_num 1 # time.sleep(0.01) g_num tmp print(f{threading.current_thread().name}: {g_num}\n, end) if __name__ __main__: g_num 0 threads [threading.Thread(targetfunc, namef线程{i}) for i in range(3)] [t.start() for t in threads] # 列表推导式,循环执行start [t.join() for t in threads] print(g_num) # 30结果为30看似没有问题这是因为这个修改操作花费的时间太短了短到我们无法想象。所以线程间轮询执行时都能获取到最新的 g_num 值。因此暴露问题的概率就变得微乎其微。我们添加0.01秒的延迟时间import time import threading def func(): global g_num for _ in range(10): tmp g_num 1 time.sleep(0.01) g_num tmp print(f{threading.current_thread().name}: {g_num}\n, end) if __name__ __main__: g_num 0 threads [threading.Thread(targetfunc, namef线程{i}) for i in range(3)] [t.start() for t in threads] [t.join() for t in threads] print(g_num) # 10对同一个数据 g_num 进行修改操作,就会遇到线程安全问题可以看到最终结果并不是30。这是因为在修改 g_num 前有0.01秒的休眠时间某个线程延时后CPU立即分配计算资源给其他线程。此时0.01秒的休眠还未结束这个线程还未将修改后的数据赋值给 g_num因此其他线程获取到的并不是最新值所以才出现上面的结果。解决办法就是 互斥锁import time import threading def func(): global g_num *# 加锁* * *lock.acquire() for _ in range(10): tmp g_num 1 time.sleep(0.01) g_num tmp print(f{threading.current_thread().name}: {g_num}\n, end) *# 释放锁* * *lock.release() if __name__ __main__: *# 创建锁* * *lock threading.Lock() g_num 0 threads [threading.Thread(targetfunc, namef线程{i}) for i in range(3)] [t.start() for t in threads] [t.join() for t in threads] print(g_num) *# 10*但是这样会导致 线程 变为同步执行,影响性能要注意锁的位置,哪里堵塞在哪里添加 锁, 放在循环内import time import threading def func(): global g_num for _ in range(10): *# 加锁* * *lock.acquire() tmp g_num 1 time.sleep(0.01) g_num tmp *# 释放锁* * *lock.release() print(f{threading.current_thread().name}: {g_num}\n, end) if __name__ __main__: *# 创建锁* * *lock threading.Lock() g_num 0 threads [threading.Thread(targetfunc, namef线程{i}) for i in range(3)] [t.start() for t in threads] [t.join() for t in threads] print(g_num) *# 10*13.5 互斥锁某个线程要更改共享数据时先将其锁定此时其他线程不能更改。直到该线程释放资源将资源的状态变成“非锁定”其他的线程才能再次锁定该资源。互斥锁保证了每次只有一个线程进行写入操作从而保证了多线程情况下数据的正确性。互斥锁的使用可以通过threading.Lock()创建互斥锁。使用lock.acquire([blockingTrue][, timeout-1])来获取锁blocking如果为True线程会阻塞直到获取到锁。如果为False线程立即返回。获取锁成功返回True否则返回False。timeout为等待的超时时间单位为秒。如果超时仍未获取到锁则返回False。。使用lock.release()释放锁。import time import threading def sale_ticket(): global ticket_num while True: *# 加锁* * *lock.acquire() if ticket_num 0: lock.release() break time.sleep(0.1) ticket_num - 1 print(f{threading.current_thread().name}卖了1张票,还剩: {ticket_num}张) *# 释放锁* * ** *lock.release() if __name__ __main__: ticket_num 100 *# 创建锁* * *lock threading.Lock() threads [threading.Thread(targetsale_ticket, namef窗口{i1}) for i in range(3)] [t.start() for t in threads] [t.join() for t in threads] print(f主线程:{ticket_num})在判断的时候也要添加 释放锁,不然无法正常退出程序13.6 GILPython 全局解释器锁Global Interpreter Lock, 简称 GIL是一个锁同一时间只允许一个线程保持 Python 解释器的控制权这意味着在任何时间点都只能有一个线程处于执行状态。执行单线程程序时看不到 GIL 的影响但它可能是 CPU 密集型和多线程代码中的性能瓶颈。GIL并不是Python的特性它是在实现Python解析器CPython时所引入的一个概念。GIL 的存在会对多线程的效率有不小影响。甚至就几乎等于Python是个单线程的程序。我们可能会想 GIL只要释放的勤快效率也不会差至少也不会比单线程的效率差。理论上是这样。但实际上Python为了让各个线程能够平均利用CPU时间会计算当前已执行的微代码数量达到一定阈值后就强制释放GIL。而这时也会触发一次操作系统的线程调度当然是否真正进行上下文切换由操作系统自主决定。从释放 GIL 到获取 GIL 之间几乎是没有间隙的。所以当其他在其他核心上的线程被唤醒时大部分情况下主线程已经又再一次获取到 GIL 了。这个时候被唤醒执行的线程只能白白的浪费CPU时间看着另一个线程拿着 GIL 执行。然后达到切换时间后进入待调度状态再被唤醒再等待以此往复恶性循环。Python的每个版本中也在逐渐改进GIL和线程调度之间的互动关系。例如先尝试持有GIL在做线程上下文切换在IO等待时释放GIL等尝试。但是无法改变的是GIL的存在使得操作系统线程调度的这个本来就昂贵的操作变得更奢侈了。总之当你的程序需要进行大量的CPU计算时GIL会成为性能的瓶颈。13.7 进程和线程对比13.7.1 区别资源分配进程拥有独立的内存空间和系统资源每个进程都有自己的代码段、数据段和堆栈等。而线程共享所属进程的内存空间和资源同一进程内的线程之间可以直接访问共享内存。开销创建进程需要分配独立的内存、打开文件等系统资源开销较大。创建线程只需在所属进程的内存空间内进行少量资源分配开销较小。并发性在多核心 CPU 环境下进程和线程都可以异步执行但进程之间的异步是真正的异步每个进程在不同核心上同时执行而线程之间的异步在单核心 CPU 上是通过时间片轮转实现的 “伪异步”在同一时刻只有一个线程执行在多核心 CPU 上可以实现异步。但是在Cpython中因为GIL的存在也不是真正的异步独立性进程之间相互独立一个进程的崩溃通常不会影响其他进程。而同一进程内的线程之间相互影响一个线程出现问题可能导致整个进程崩溃。通信进程间通信相对复杂需要使用特殊的机制如管道、消息队列、共享内存等。线程间通信相对简单因为它们共享内存可以直接访问共享变量。13.7.2 使用场景适合使用多线程的情况**I/O 密集型任务**如网络请求、文件读写等。线程共享内存切换开销小在等待 I/O 操作完成的时间内可以切换到其他线程执行提高整体效率。例如一个程序需要同时从多个网站下载数据使用多线程可以在等待网络响应时执行其他下载任务。**对资源共享要求高**线程间共享内存方便数据共享和通信。例如在一个图形界面程序中多个线程需要共享界面数据并进行实时更新。适合使用多进程的情况CPU 密集型任务多进程可以利用多核心 CPU 实现真正的并行计算充分发挥硬件性能。例如进行复杂的科学计算、数据处理等任务每个进程在不同核心上独立计算提高计算速度。**需要隔离的任务**进程相互独立一个进程崩溃不会影响其他进程。对于一些可能出现异常或不稳定的任务使用多进程可以保证系统的稳定性。例如运行多个独立的服务每个服务作为一个进程避免一个服务出错影响其他服务。