进程 (process)
进程是对各种资源管理的集合,包含对各种资源的调用、内存的管理、网络接口的调用 进程要操作 CPU 必须先启动一个线程,启动一个进程的时候会自动创建一个线程,进程里的第一个线程就是主线程 程序执行的实例 有唯一的进程标识符(pid)
multiprossing 模块
启动进程
示例:
1import multiprocessing 2import time 3 4 5def process_run(n): 6 time.sleep(1) 7 print('process', n) 8 9 10for i in range(10): 11 p = multiprocessing.Process(target=process_run, args=(i, )) 12 p.start() 13
所有进程都是由父进程启动的
示例:
1import multiprocessing 2import os 3 4 5def show_info(title): 6 print(title) 7 print('module name:', __name__) 8 print('parent process:', os.getppid()) 9 print('process id', os.getpid()) 10 print('\n\n') 11 12 13def f(name): 14 show_info('function f') 15 print(name) 16 17 18if __name__ == '__main__': 19 show_info('main process line') 20 p = multiprocessing.Process(target=f, args=('children process', )) 21 p.start() 22
进程间通信
线程间共享内存空间,进程间只能通过其他方法进行通信
Queue
注意这个 Queue 不同于 queue.Queue Queue type using a pipe, buffer and thread 两个进程的 Queue 并不是同一个,而是将数据 pickle 后传给另一个进程的 Queue 用于父进程与子进程之间的通信或同一父进程的子进程之间通信
示例:
1from multiprocessing import Process, Queue 2 3 4def p_put(*args): 5 q.put(args) 6 print('Has put %s' % args) 7 8 9def p_get(*args): 10 print('%s wait to get...' % args) 11 print(q.get()) 12 print('%s got it' % args) 13 14 15q = Queue() 16p1 = Process(target=p_put, args=('p1', )) 17p2 = Process(target=p_get, args=('p2', )) 18p1.start() 19p2.start() 20
输出结果: Has put p1 p2 wait to get... ('p1',) p2 got it
换成 queue 示例:
1from multiprocessing import Process 2import queue 3 4 5def p_put(*args): 6 q.put(args) 7 print('Has put %s' % args) 8 9 10def p_get(*args): 11 print('%s wait to get...' % args) 12 print(q.get()) 13 print('%s got it' % args) 14 15 16q = queue.Queue() 17p1 = Process(target=p_put, args=('p1', )) 18p2 = Process(target=p_get, args=('p2', )) 19p1.start() 20p2.start() 21
输出结果: Has put p1 p2 wait to get...
由于父进程启动子进程时是复制一份,所以每个子进程里也有一个空的队列,但是这些队列数据独立,所以 get 时会阻塞
Pipe
Pipe(管道) 是通过 socket 进行进程间通信的 所以步骤与建立 socket 连接相似: 建立连接、发送/接收数据(一端发送另一端不接受就会阻塞)、关闭连接 示例:
1from multiprocessing import Pipe, Process 2 3 4def f(conn): 5 conn.send('send by child') 6 print('child recv:', conn.recv()) 7 conn.close() 8 9 10parent_conn, child_conn = Pipe() # 获得 Pipe 连接的两端 11p = Process(target=f, args=(child_conn, )) 12p.start() 13print('parent recv:', parent_conn.recv()) 14parent_conn.send('send by parent') 15p.join() 16
输出结果: parent recv: send by child child recv: send by parent
进程间数据共享
Manager
Manager 实现的是进程间共享数据 支持的可共享数据类型:
1list 2dict 3Value 4Array 5Namespace 6Queue queue.Queue 7JoinableQueue queue.Queue 8Event threading.Event 9Lock threading.Lock 10RLock threading.RLock 11Semaphore threading.Semaphore 12BoundedSemaphore threading.BoundedSemaphore 13Condition threading.Condition 14Barrier threading.Barrier 15Pool pool.Pool
示例:
1from multiprocessing import Manager, Process 2import os 3 4 5def func(): 6 m_dict['key'] = 'value' 7 m_list.append(os.getpid()) 8 9 10manager = Manager() 11m_dict = manager.dict() 12m_list = manager.list() 13p_list = [] 14for i in range(10): 15 p = Process(target=func) 16 p.start() 17 p_list.append(p) 18for p in p_list: 19 p.join() 20print(m_list) 21print(m_dict) 22
进程锁
打印时可能会出错,加锁可以避免 示例:
1from multiprocessing import Lock, Process 2 3 4def foo(n, l): 5 l.acquire() 6 print('hello world', n) 7 l.release() 8 9 10lock = Lock() 11for i in range(100): 12 process = Process(target=foo, args=(i, lock)) 13 process.start() 14
进程池 (pool)
同一时间最多有几个进程在 CPU 上运行 示例:
1from multiprocessing import Pool 2import time 3import os 4 5 6def foo(n): 7 time.sleep(1) 8 print('In process', n, os.getpid()) 9 return n 10 11 12def bar(*args): 13 print('>>done: ', args, os.getpid()) 14 15 16pool = Pool(processes=3) 17print('主进程: ', os.getpid()) 18for i in range(10): 19 # pool.apply(func=foo, args=(i, )) 20 pool.apply_async(func=foo, args=(i, ), callback=bar) 21print('end') 22pool.close() 23pool.join() 24
从程序运行过程中可以看出:同一时间最多只有3个进程在运行,类似于线程中的信号量 主进程在执行 callback 函数
注意
1. pool.apply(func=foo, args=(i, )) 是串行执行 pool.apply_async(func=foo, args=(i, ), callback=bar) 是并行执行 2. callback 函数会以 target 函数返回结果为参数,在 target 函数执行结束之后执行 callback 函数是主进程调用的 3. 如果不执行 join,程序会在主进程执行完成之后直接结束,不会等待子进程执行完成 Pool.join() 必须在 Pool.close() 之后执行,否则会报错:ValueError: Pool is still running