YAOTU INSIGHTS

python threading和multiprocessing模块基本用法实例分析

python threading和multiprocessing模块基本用法实例分析
现在本文给大家详细讲讲这个情况, 就是关于那个模块的基本用法怎么操作。下面我会把它分享出来, 大家看一看可以参考一下, 具体的内容如下所示:前言这几天, 为了做一个小项目, 我研究了一下并发编程。所谓并发, 无非就是多线程和多进程。最初找到的模块是那个, 因为我的印象中认为线程是“轻量...”, “切换快...”, “可共享进程资源...”等等。但是, 我没有想到这里面的水深不见底, 进而找到了更好的替代品是另一个模块。下面会讲一些在使用其中的经验。在后面所展示的种种代码方面, 经过测试能够确认是全部通过的, 而测试所处的环境呢, 具体来说是.04加上.6.5这样的版本组合。一、使用模块创建线程1、三种线程创建方式1传入一个函数这种做法是最为普通的, 就是去调用当前那个类的构造函数, 并且给参数赋予等于func的值, 然后再去使用由构造函数所返回的那个实例来对这个叫做start的东西进行调用, 也就是让该线程开始运行起来, 这个被启动的线程将会去执行这个函数func, 当然了, 要是说这个函数func它是需要参数的话, 那就完全可以在上述的构造函数里面直接把这个名为args的参数传进去, 其形式类似于将括号包裹起来的若干个值放在符号的右边, 下面展示一下相关的示例代码片段, 以便于理解:#!/usr/bin/python#-*-coding:utf-8-*-import threading#用于线程执行的函数def counter(n):cnt 0;for i in xrange(n):for j in xrange(i):cnt j;print cnt;if __name__ __main__:#初始化一个线程对象传入函数counter及其参数1000th threading.Thread(targetcounter, args(1000,));#启动线程th.start();#主线程阻塞等待子线程结束th.join();这段代码看起来非常直观, 其中的函数其实就是一个结构很平凡的两层循环操作。不过有一点地方特别需要大家去重视, 那就是 th.join()这个方法的调用。它的具体功能是指让当前运行的主线程处于一种自我暂停的状态, 并且只有当由 th 变量所代表的那个分线程真正执行完毕之后, 整体程序才会继续走下面的流程。假如我们在代码里删掉了这一句内容的话, 那么整个运行过程就会在触发后的瞬间戛然而止。关于 join 一词背后的含义确实显得有些晦涩难懂, 但实际上我们完全可以把它想象成类似于 “while th.( ) : time.sleep(1)”这中一种表现形式, 这样来理解的话, 想必会变得更容易掌握一些。尽管二者的意思是一样的, 不过, 在后面我们将看到的事实是, 使用join这种方法存在着一些陷阱。2传入一个可调用的对象许多对象都是我们所说的可调用的, 也就是说这些对象是任何能够通过函数操作符小括号来调用的对象, 大家可以在《核心编程》第14章里看到这一点。类的对象也是可以直接被调用的, 当这些类对象被调用时, 系统会自动去调用这个对象内部的内置方法小括号所以这种新建线程的方法其实就是给线程指定一个把方法进行了重载的具体的对象。示例代码如下:#!/usr/bin/python#-*-coding:utf-8-*-import threading#可调用的类class Callable(object):def __init__(self, func, args):self.func func;self.args args;def __call__(self):apply(self.func, self.args);#用于线程执行的函数def counter(n):cnt 0;for i in xrange(n):for j in xrange(i):cnt j;print cnt;if __name__ __main__:#初始化一个线程对象传入可调用的Callable对象并用函数counter及其参数1000初始化这个对象th threading.Thread(targetCallable(counter, (1000,)));#启动线程th.start();#主线程阻塞等待子线程结束th.join();在这个例子中, 最关键的一句话是apply(self.func, self.args);。在这一部分代码里, 我们使用在初始化阶段所传入的那个函数对象, 以及同时传入的参数列表, 来执行一次调用操作。3继承类这种方式是利用继承这个类, 并且把这个run方法给重写覆盖掉的手段, 去实现出自己定制的线程操作功能的, 下面这些代码内容就是具体的实例展现:#!/usr/bin/python#-*-coding:utf-8-*-import threading, time, randomdef counter():cnt 0;for i in xrange(10000):for j in xrange(i):cnt j;class SubThread(threading.Thread):def __init__(self, name):threading.Thread.__init__(self, namename);def run(self):i 0;while i 4:print self.name,counting...\n;counter();print self.name,finish\n;i 1;if __name__ __main__:th SubThread(thread-1);th.start();th.join();print all done;这个例子定义了一个类, 它从另一个类继承而来, 而且重写了 run 方法, 在这个 run 方法里面去调用原有的逻辑, 并且打印一些相关的信息, 大家可以看到这种方式是非常直观、容易理解的, 不过需要注意的是, 在构造函数当中, 一定不要忘记先去调用父类的构造函数来完成必要的初始化的操作。2、多线程的限制多线程存在一个令人感到困扰的限制, 也就是全局解释器锁 lock。这个锁的意思是说, 在任一特定的时间点上, 只能允许一个线程来使用解释器。这种情况和单个中央处理器执行多个程序运行的模式是一样的。所有线程都是轮流获取资源去使用的。这种现象被称作是“并发”。它并不是真正意义上的“并行”。官方技术手册里的说明指出, 这样设计的目的是为了保障对象模型的正确性得以实现。由这个锁带来的麻烦, 如果有一个属于计算密集型的线程把cpu给占住了, 剩下的其它线程都只能干等着……你试着去想一下, 在你手头的那多个线程里面, 要是出现了这么一种线程, 情况是多么让人沮丧, 原本的多线程运行模式, 硬生生就变成了一种串行的处理方式不过呢, 这个模块也并非是没有一点作用的, 技术手册上面也给出了说明, 也就是当被用来对待那些属于IO密集型的任务之时在IO发生进行的整个过程里面, 当前的线程会主动释放掉解释器的持有权, 正因为这样, 另外的那些线程才可以获取到使用解释器的一个机会。因此在决定要不要使用这个模块的时候, 是需要考虑到所要面对的具体的任务类型是怎么样的。二、使用创建进程1、三种创建方式进程的创建方法跟线程是完全没有区别的, 唯一需要做的就是把中间的点号换成点号就行, 相关的模块在尽力保持和那个模块的名字完全一致的情况下去进行操作, 如果想要看具体的例子, 就去看一看上面讲到的关于线程部分的参考代码, 我们这里之所以只给出第一种通过函数来进行操作的用法, 是因为其他的方式可以参考之前的内容。#!/usr/bin/python#-*-coding:utf-8-*-import multiprocessing, timedef run():i 0;while i10000:print running;time.sleep(2);i 1;if __name__ __main__:p multiprocessing.Process(targetrun);p.start();#p.join();print p.pid;print master gone;2、创建进程池这个模块还支持先一次性创建多个进程, 然后才对它们分配具体任务。关于这部分的具体内容, 大家可以去仔细查看相关的手册文档, 毕竟我对这一块的研究投入精力不多, 所以不敢妄加猜测或编写不准确的信息。pool multiprocessing.Pool(processes4)pool.apply_async(func, args...)3、使用进程的好处它是完全并行的。并且没有GIL的限制。所以能够利用多CPU多核的环境。它还可以接受Linux信号。在后面会看到。这个功能非常的好用。三、实例研究在这个假设的情景里面, 我们的主要任务是, 有一个主进程要去启动好多个子进程, 让这些子进程各自去处理不一样的任务, 而且这些子进程可能还会生成它们自己的线程, 用这些线程来做一些输入输出的事情, 因为前面已经讲过了, 在执行输入输出操作的时候, 线程的表现是挺不错的, 现在需要实现的功能内容是, 我们要向这些子进程发送信号, 要让这些信号得到正确的响应, 比如说, 一旦发出了信号, 子进程就应该去通知它里面的线程去完成工作, 然后就要以一种不那么粗暴的方式结束运行。目前需要对付的难题是, 第一点是在子类化的对象里面怎么抓住信号第二点是怎么样做到“优雅的退出”。接下来我们要一个一个来说清楚。1、子类化并捕捉信号如果在采用的第一种进程创建方式里面, 也就是把函数当作参数送进去的这种做法, 捕捉信号就变得十分容易了。不妨假设给这个由进程去执行的函数取一个名字叫做func, 然后代码示例的部分可以参照如下所示:#!/usr/bin/python#-*-coding:utf-8-*-import multiprocessing, signal,timedef handler(signum, frame):print signal, signum;def run():signal.signal(signal.SIGTERM, handler);signal.signal(signal.SIGINT, handler);i 0;while i10000:print running;time.sleep(2);i 1;if __name__ __main__:p multiprocessing.Process(targetrun);p.start();#p.join();print p.pid;print master gone;这段代码是在第一种创建方式的基础之上修改而来的, 增加了两行调用代码, 具体的写法是加上括号然后省略号, 这一做法的意思就是说这个函数需要去捕捉两个特定的信号, 除此之外, 还额外增加了一个新的函数, 这个新函数的主要用途就是在捕捉到对应的信号之时进行相应的处理操作, 在当前这种情况下, 我们只是做了非常简单的一个动作, 也就是把捕捉到的信号的数值打印出来。请注意, p.join()这个方法是被注释掉了的。这和之前的线程那个情况有一点不太一样。因为新的进程在启动以后就开始运行了。主进程也不需要去等它运行完。该干什么就继续干什么去吧。这段代码跑起来之后会打印出子进程的进程id。根据这个拿到的id。你在别的终端输入kill -TERM id。会发现刚才的那个终端又打印出来了数字15。然而, 采用传入函数这种做法存在一个不足之处, 那就是其封装性表现欠佳。倘若功能变得相对复杂一些, 那么就必然会有很多全局变量暴露在外部环境中。最为妥当的办法, 应当是将相关功能打包封装成一个类。既然如此, 在使用这个类的时候, 又该如何去注册处理信号对应的函数呢?之前的那些示例案例, 似乎仅仅支持使用单一的、作为全局存在的函数来处理。相关的官方手册资料之中, 也没有展现出如何在类内部来处理信号的实例信息。实际上, 其解决思路与先前的办法是相似的, 过程也相当简单易懂。之前看到的那篇帖子, 为我提供了重要的灵感指引。class Master(multiprocessing.Process):def __init__(self):super(Master,self).__init__();signal.signal(signal.SIGTERM, self.handler); #注册信号处理函数self.live 1;#信号处理函数def handler(self, signum, frame):print signal:,signum;self.live 0;def run(self):print PID:,self.pid;while self.live:print living...time.sleep(2);这个方法是相当直观的一一先是在构造函数里把信号处理函数给它注册上, 接着是定义一个办法, 让它来做那个处理的活儿。那么这个进程类, 就会每隔两秒钟就打印出一个点点dot, 等它接收到那个东西之后, 会把self.live的值改成别的。然后, run方法的循环如果发现那个值是零了, 那就说明结束了, 进程也就跟着结束了。2、让进程优雅的退出现在将这次设想任务的完整代码予以呈现。我是在主进程内部启动了一个子进程, 具体做法是对某个类进行子类化来实现这一操作。接着, 该子进程被创建出来后, 又会进一步产生两个子线程。这两个子线程的用途是用来模拟所谓的生产者与消费者这一模型。在这对线程之间, 它们借由一个队列来进行信息的交流。因为存在需要互相排斥地访问这个共享队列的情形, 所以理所当然地必须添加一把锁。这把锁与普通的 Lock 对象大体相似, 但是额外具备等待以及通知这两项功能。对于生产者而言, 它每一次的操作是制造出一个随机数, 并将其投放入队列之中, 随后进入一段持续时间也是随机确定的休息状态。而消费者则是执行从队列里提取出每一个数值的动作。至于处于子进程内部的主线程, 其核心职责在于接收外部的信号指令, 以此确保整个运行流程能够实现优雅且平滑的结束。代码如下#!/usr/bin/python#-*-coding:utf-8-*-import time, multiprocessing, signal, threading, random, time, Queueclass Master(multiprocessing.Process):def __init__(self):super(Master,self).__init__();signal.signal(signal.SIGTERM, self.handler);#这个变量要传入线程用于控制线程运行为什么用dict充分利用线程间共享资源的特点#因为可变对象按引用传递标量是传值的不信写成self.live true试试self.live {stat:True};def handler(self, signum, frame):print signal:,signum;self.live[stat] 0; #置这个变量为0通知子线程可以“收工”了def run(self):print PID:,self.pid;cond threading.Condition(threading.Lock()); #创建一个condition对象用于子线程交互q Queue.Queue(); #一个队列sender Sender(cond, self.live, q); #传入共享资源geter Geter(cond, self.live, q);sender.start(); #启动线程geter.start();signal.pause(); #主线程睡眠并等待信号while threading.activeCount()-1: #主线程收到信号并被唤醒后检查还有多少线程活着除掉自己time.sleep(2); #再睡眠等待确保子线程都安全的结束print checking live, threading.activeCount();print mater gone;class Sender(threading.Thread):def __init__(self, cond, live, queue):super(Sender, self).__init__(namesender);self.cond cond;self.queue queue;self.live livedef run(self):cond self.cond;while self.live[stat]: #检查这个进程内的“全局”变量为真就继续运行cond.acquire(); #获得锁以便控制队列i random.randint(0,100);self.queue.put(i,False);if not self.queue.full():print sender add:,i;cond.notify(); #唤醒等待锁的其他线程cond.release(); #释放锁time.sleep(random.randint(1,3));print sender doneclass Geter(threading.Thread):def __init__(self, cond, live, queue):super(Geter, self).__init__(namegeter);self.cond cond;self.queue queue;self.live livedef run(self):cond self.cond;while self.live[stat]:cond.acquire();if not self.queue.empty():i self.queue.get();print geter get:,i;cond.wait(3);cond.release();time.sleep(random.randint(1,3));print geter doneif __name__ __main__:master Master();master.start(); #启动子进程必须要留意的一个细节在于, 在咱们的run方法里面, 紧跟着.start()和geter.start()这两个动作之后, 按照通常的逻辑, 紧接着就应该去调用.join()和geter.join()了, 这么做是为了让主线程去等待子线程执行结束, 之前提到过的.join()方法的陷阱正是出现于此, 因为.join()这个方法会把主线程给阻塞住, 导致主线程没有办法再继续来捕捉信号, 所以在刚开始钻研这块内容的时候, 一度以为是自己写的信号处理函数出现了错误。在网上展开讨论的数量确实是比较少的, 所以这里将相关内容解释清楚了。参考《核心编程》《 》更多关于相关内容感兴趣的读者可查看本站专题《进程与线程操作技巧总结》、《数据结构与算法教程》、《函数使用技巧总结》、《字符串操作技巧汇总》、《入门与进阶经典教程》、《MySQL数据库程序设计入门教程》及《常见数据库操作技巧汇总》本文所提及到的内容, 是希望可以给大家在程序设计方面提供一些相应的帮助。