1、前述
进程 线程 协程 异步
并发编程(不是并行)目前有四种方式:多进程、多线程、协程和异步。
- 多进程编程在python中有类似C的os.fork,更高层封装的有multiprocessing标准库
- 多线程编程python中有Thread和threading
- 异步编程在linux下主+要有三种实现select,poll,epoll
- 协程在python中通常会说到yield,关于协程的库主要有greenlet,stackless,gevent,eventlet等实现
进程
- 不共享任何状态
- 调度由操作系统完成
- 有独立的内存空间(上下文切换的时候需要保存栈、cpu寄存器、虚拟内存、以及打开的相关句柄等信息,开销大)
- 通讯主要通过信号传递的方式来实现(实现方式有多种,信号量、管道、事件等,通讯都需要过内核,效率低)
线程
- 共享变量(解决了通讯麻烦的问题,但是对于变量的访问需要加锁)
- 调度由操作系统完成(由于共享内存,上下文切换变得高效)
- 一个进程可以有多个线程,每个线程会共享父进程的资源(创建线程开销占用比进程小很多,可创建的数量也会很多)
- 通讯除了可使用进程间通讯的方式,还可以通过共享内存的方式进行通信(通过共享内存通信比通过内核要快很多)
协程
- 调度完全由用户控制
- 一个线程(进程)可以有多个协程
- 每个线程(进程)循环按照指定的任务清单顺序完成不同的任务(当任务被堵塞时,执行下一个任务;当恢复时,再回来执行这个任务;任务间切换只需要保存任务的上下文,没有内核的开销,可以不加锁的访问全局变量)
- 协程需要保证是非堵塞的且没有相互依赖
- 协程基本上不能同步通讯,多采用异步的消息通讯,效率比较高
总结
1协程 > 多进程 >多线程 2使用gevent,可以获得极高的并发性能,但gevent只能在Unix/Linux下运行,在Windows下不保证正常安装和运行
既然有了线程为什么还要协程呢?因为线程是系统级别的,在做切换的时候消耗是特别大的,具体为什么这么大等我研究好了再告诉你;同时线程的切换是由CPU决定的,可能你刚好执行到一个地方的时候就要被迫终止,这个时候你需要用各种措施来保证你的数据不出错,所以线程对于数据安全的操作是比较复杂的。而协程是用户级别的切换,且切换是由自己控制,不受外力终止.
- 进程拥有自己独立的堆和栈,既不共享堆,亦不共享栈,进程由操作系统调度
- 线程拥有自己独立的栈和共享的堆,共享堆,不共享栈,线程亦由操作系统调度(标准线程是的)
- 协程和线程一样共享堆,不共享栈,协程由程序员在协程的代码里显示调度
聊聊协程
协程,又称微线程,纤程。 Python的线程并不是标准线程,是系统级进程,线程间上下文切换有开销,而且Python在执行多线程时默认加了一个全局解释器锁(GIL),因此Python的多线程其实是串行的,所以并不能利用多核的优势,也就是说一个进程内的多个线程只能使用一个CPU。
1def coroutine(func): 2 def ret(): 3 f = func() 4 f.next() 5 return f 6 return ret 7 8 9@coroutine 10def consumer(): 11 print "Wait to getting a task" 12 while True: 13 n = (yield) 14 print "Got %s",n 15 16 17import time 18def producer(): 19 c = consumer() 20 task_id = 0 21 while True: 22 time.sleep(1) 23 print "Send a task to consumer" % task_id 24 c.send("task %s" % task_id) 25 26if __name__ == "__main__": 27 producer() 28 29运行结果: 30 31Wait to getting a task 32Send a task 0 to consumer 33Got task 0 34Send a task 1 to consumer 35Got task 1 36Send a task 2 to consumer 37Got task 2
传统的生产者-消费者模型是一个线程写消息,一个线程取消息,通过锁机制控制队列和等待,但容易死锁。 如果改用协程,生产者生产消息后,直接通过yield跳转到消费者开始执行,待消费者执行完毕后,切换回生产者继续生产,效率极高。
2、Gevent
介绍
gevent是基于协程的Python网络库。特点:
- 基于libev的快速事件循环(Linux上epoll,FreeBSD上kqueue)。
- 基于greenlet的轻量级执行单元。
- API的概念和Python标准库一致(如事件,队列)。
- 可以配合socket,ssl模块使用。
- 能够使用标准库和第三方模块创建标准的阻塞套接字(gevent.monkey)。
- 默认通过线程池进行DNS查询,也可通过c-are(通过GEVENT_RESOLVER=ares环境变量开启)。
- TCP/UDP/HTTP服务器
- 子进程支持(通过gevent.subprocess)
- 线程池
核心部分
- Greenlets
- 同步和异步执行
- 确定性
- 创建Greenlets
- Greenlet状态
- 程序停止
- 超时
- 猴子补丁
Greenlets
gevent中的主要模式, 它是以C扩展模块形式接入Python的轻量级协程。 全部运行在主程序操作系统进程的内部,但它们被程序员协作式地调度。
在任何时刻,只有一个协程在运行。
区别于multiprocessing、threading等提供真正并行构造的库, 这些库轮转使用操作系统调度的进程和线程,是真正的并行。
同步和异步执行
并发的核心思想在于,大的任务可以分解成一系列的子任务,后者可以被调度成 同时执行或异步执行,而不是一次一个地或者同步地执行。两个子任务之间的 切换也就是上下文切换。
在gevent里面,上下文切换是通过yielding来完成的.
1import gevent 2def foo(): 3 print('Running in foo') 4 gevent.sleep(0) 5 print('Explicit context switch to foo again') 6def bar(): 7 print('Explicit context to bar') 8 gevent.sleep(0) 9 print('Implicit context switch back to bar') 10gevent.joinall([ 11 gevent.spawn(foo), 12 gevent.spawn(bar), 13]) 14 15Running in foo 16Explicit context to bar 17Explicit context switch to foo again 18Implicit context switch back to bar
网络延迟或IO阻塞隐式交出greenlet上下文的执行权.
1import time 2import gevent 3from gevent import select 4start = time.time() 5tic = lambda: 'at %1.1f seconds' % (time.time() - start) 6def gr1(): 7 print('Started Polling: %s' % tic()) 8 select.select([], [], [], 1) 9 print('Ended Polling: %s' % tic()) 10def gr2(): 11 print('Started Polling: %s' % tic()) 12 select.select([], [], [], 2) 13 print('Ended Polling: %s' % tic()) 14def gr3(): 15 print("Hey lets do some stuff while the greenlets poll, %s" % tic()) 16 gevent.sleep(1) 17gevent.joinall([ 18 gevent.spawn(gr1), 19 gevent.spawn(gr2), 20 gevent.spawn(gr3), 21]) 22 23运行结果: 24Started Polling: at 0.0 seconds 25Started Polling: at 0.0 seconds 26Hey lets do some stuff while the greenlets poll, at 0.0 seconds 27Ended Polling: at 1.0 seconds 28Ended Polling: at 2.0 seconds
同步vs异步
1import gevent 2import random 3def task(pid): 4 gevent.sleep(random.randint(0,2)*0.001) 5 print('Task %s done' % pid) 6def synchronous(): 7 for i in xrange(5): 8 task(i) 9def asynchronous(): 10 threads = [gevent.spawn(task, i) for i in xrange(5)] 11 gevent.joinall(threads) 12 print('Synchronous:') 13 synchronous() 14 print('Asynchronous:') 15 asynchronous() 16 17运行结果: 18 19Synchronous: 20Task 0 done 21Task 1 done 22Task 2 done 23Task 3 done 24Task 4 done 25Asynchronous: 26Task 2 done 27Task 0 done 28Task 1 done 29Task 3 done 30Task 4 done
确定性
greenlet具有确定性。在相同配置相同输入的情况下,它们总是会产生相同的输出
1import time 2def echo(i): 3 time.sleep(0.001) 4 return i 5# Non Deterministic Process Pool 6from multiprocessing.pool import Pool 7p = Pool(10) 8run1 = [a for a in p.imap_unordered(echo, xrange(10))] 9run2 = [a for a in p.imap_unordered(echo, xrange(10))] 10run3 = [a for a in p.imap_unordered(echo, xrange(10))] 11run4 = [a for a in p.imap_unordered(echo, xrange(10))] 12print(run1 == run2 == run3 == run4) 13# Deterministic Gevent Pool 14from gevent.pool import Pool 15p = Pool(10) 16run1 = [a for a in p.imap_unordered(echo, xrange(10))] 17run2 = [a for a in p.imap_unordered(echo, xrange(10))] 18run3 = [a for a in p.imap_unordered(echo, xrange(10))] 19run4 = [a for a in p.imap_unordered(echo, xrange(10))] 20print(run1 == run2 == run3 == run4) 21 22运行结果: 23 24False 25True
即使gevent通常带有确定性,当开始与如socket或文件等外部服务交互时, 不确定性也可能溜进你的程序中。因此尽管gevent线程是一种“确定的并发”形式, 使用它仍然可能会遇到像使用POSIX线程或进程时遇到的那些问题。
涉及并发长期存在的问题就是竞争条件(race condition)(当两个并发线程/进程都依赖于某个共享资源同时都尝试去修改它的时候, 就会出现竞争条件),这会导致资源修改的结果状态依赖于时间和执行顺序。 这个问题,会导致整个程序行为变得不确定。
解决办法: 始终避免所有全局的状态.
创建Greenlets
gevent对Greenlet初始化提供了一些封装.
1import gevent 2from gevent import Greenlet 3def foo(message, n): 4 gevent.sleep(n) 5 print(message) 6 thread1 = Greenlet.spawn(foo, "Hello", 1) 7 thread2 = gevent.spawn(foo, "I live!", 2) 8 thread3 = gevent.spawn(lambda x: (x+1), 2) 9 threads = [thread1, thread2, thread3] 10 gevent.joinall(threads) 11 12执行结果: 13 14Hello 15I live!
除使用基本的Greenlet类之外,你也可以子类化Greenlet类,重载它的_run方法.
1import gevent 2from gevent import Greenlet 3class MyGreenlet(Greenlet): 4 def __init__(self, message, n): 5 Greenlet.__init__(self) 6 self.message = message 7 self.n = n 8 def _run(self): 9 print(self.message) 10 gevent.sleep(self.n) 11g = MyGreenlet("Hi there!", 3) 12g.start() 13g.join() 14 15执行结果: 16Hi there!
Greenlet状态
greenlet的状态通常是一个依赖于时间的参数:
- started – Boolean, 指示此Greenlet是否已经启动
- ready() – Boolean, 指示此Greenlet是否已经停止
- successful() – Boolean, 指示此Greenlet是否已经停止而且没抛异常
- value – 任意值, 此Greenlet代码返回的值
- exception – 异常, 此Greenlet内抛出的未捕获异常
程序停止
当主程序(main program)收到一个SIGQUIT信号时,不能成功做yield操作的 Greenlet可能会令意外地挂起程序的执行。这导致了所谓的僵尸进程, 它需要在Python解释器之外被kill掉
通用的处理模式就是在主程序中监听SIGQUIT信号,调用gevent.shutdown退出程序
1import gevent 2import signal 3def run_forever(): 4 gevent.sleep(1000) 5 if __name__ == '__main__': 6 gevent.signal(signal.SIGQUIT, gevent.shutdown) 7 thread = gevent.spawn(run_forever) 8 thread.join()
超时
通过超时可以对代码块儿或一个Greenlet的运行时间进行约束
1import gevent 2from gevent import Timeout 3seconds = 10 4timeout = Timeout(seconds) 5timeout.start() 6def wait(): 7 gevent.sleep(10) 8 try: 9 gevent.spawn(wait).join() 10 except Timeout: 11 print('Could not complete')
超时类
1import gevent 2from gevent import Timeout 3time_to_wait = 5 # seconds 4 class TooLong(Exception): 5 pass 6 with Timeout(time_to_wait, TooLong): 7 gevent.sleep(10)
另外,对各种Greenlet和数据结构相关的调用,gevent也提供了超时参数
1import gevent 2from gevent import Timeout 3def wait(): 4 gevent.sleep(2) 5timer = Timeout(1).start() 6thread1 = gevent.spawn(wait) 7try: 8 thread1.join(timeout=timer) 9except Timeout: 10 print('Thread 1 timed out') 11# -- 12timer = Timeout.start_new(1) 13thread2 = gevent.spawn(wait) 14try: 15 thread2.get(timeout=timer) 16except Timeout: 17 print('Thread 2 timed out') 18# -- 19try: 20 gevent.with_timeout(1, wait) 21except Timeout: 22 print('Thread 3 timed out') 23 24运行结果: 25 26Thread 1 timed out 27Thread 2 timed out 28Thread 3 timed out
猴子补丁(Monkey patching)
gevent的死角.
1import socket 2print(socket.socket) 3print("After monkey patch") 4from gevent import monkey 5monkey.patch_socket() 6print(socket.socket) 7import select 8print(select.select) 9monkey.patch_select() 10print("After monkey patch") 11print(select.select) 12 13运行结果: 14 15class 'socket.socket' 16After monkey patch 17class 'gevent.socket.socket' 18built-in function select 19After monkey patch 20function select at 0x1924de8
Python的运行环境允许我们在运行时修改大部分的对象,包括模块,类甚至函数。 这是个一般说来令人惊奇的坏主意,因为它创造了“隐式的副作用”,如果出现问题 它很多时候是极难调试的。虽然如此,在极端情况下当一个库需要修改Python本身 的基础行为的时候,猴子补丁就派上用场了。在这种情况下,gevent能够修改标准库里面大部分的阻塞式系统调用,包括socket、ssl、threading和 select等模块,而变为协作式运行。
例如,Redis的python绑定一般使用常规的tcp socket来与redis-server实例通信。 通过简单地调用gevent.monkey.patch_all(),可以使得redis的绑定协作式的调度 请求,与gevent栈的其它部分一起工作。
这让我们可以将一般不能与gevent共同工作的库结合起来,而不用写哪怕一行代码。 虽然猴子补丁仍然是邪恶的(evil),但在这种情况下它是“有用的邪恶(useful evil)”
数据结构
- 事件(event)
- 队列(Queue)
- 组和池(group/pool)
- 锁和信号量
- 线程局部变量
- 子进程
- Actors
事件(event)
事件(event)是一个在Greenlet之间异步通信的形式
1import gevent 2from gevent.event import Event 3 4evt = Event() 5 6def setter(): 7 print('A: Hey wait for me, I have to do something') 8 gevent.sleep(3) 9 print("Ok, I'm done") 10 evt.set() 11def waiter(): 12 print("I'll wait for you") 13 evt.wait() # blocking 14 print("It's about time") 15def main(): 16 gevent.joinall([ 17 gevent.spawn(setter), 18 gevent.spawn(waiter), 19 gevent.spawn(waiter), 20 gevent.spawn(waiter) 21 ]) 22if __name__ == '__main__': 23 main() 24 25运行结果: 26 27A: Hey wait for me, I have to do something 28I'll wait for you 29I'll wait for you 30I'll wait for you 31Ok, I'm done 32It's about time 33It's about time 34It's about time
事件对象的一个扩展是AsyncResult,它允许你在唤醒调用上附加一个值。 它有时也被称作是future或defered,因为它持有一个指向将来任意时间可设置为任何值的引用
1import gevent 2from gevent.event import AsyncResult 3a = AsyncResult() 4def setter(): 5 gevent.sleep(3) 6 a.set('Hello!') 7def waiter(): 8 print(a.get()) 9gevent.joinall([ 10 gevent.spawn(setter), 11 gevent.spawn(waiter), 12])
队列(Queue)
队列是一个排序的数据集合,它有常见的put / get操作, 但是它是以在Greenlet之间可以安全操作的方式来实现的
1import gevent 2from gevent.queue import Queue 3tasks = Queue() 4def worker(n): 5 while not tasks.empty(): 6 task = tasks.get() 7 print('Worker %s got task %s' % (n, task)) 8 gevent.sleep(0) 9 print('Quitting time!') 10def boss(): 11 for i in xrange(1,10): 12 tasks.put_nowait(i) 13gevent.spawn(boss).join() 14gevent.joinall([ 15 gevent.spawn(worker, 'steve'), 16 gevent.spawn(worker, 'john'), 17 gevent.spawn(worker, 'nancy'), 18]) 19 20执行结果: 21 22Worker steve got task 1 23Worker john got task 2 24Worker nancy got task 3 25Worker steve got task 4 26Worker john got task 5 27Worker nancy got task 6 28Worker steve got task 7 29Worker john got task 8 30Worker nancy got task 9 31Quitting time! 32Quitting time! 33Quitting time!
put和get操作都是阻塞的,put_nowait和get_nowait不会阻塞, 然而在操作不能完成时抛出gevent.queue.Empty或gevent.queue.Full异常
组和池(group/pool)
组(group)是一个运行中greenlet集合,集合中的greenlet像一个组一样会被共同管理和调度。 它也兼饰了像Python的multiprocessing库那样的平行调度器的角色,主要用在在管理异步任务的时候进行分组.
1import gevent 2from gevent.pool import Group 3def talk(msg): 4 for i in xrange(2): 5 print(msg) 6g1 = gevent.spawn(talk, 'bar') 7g2 = gevent.spawn(talk, 'foo') 8g3 = gevent.spawn(talk, 'fizz') 9group = Group() 10group.add(g1) 11group.add(g2) 12group.join() 13group.add(g3) 14group.join() 15 16运行结果: 17 18bar 19bar 20foo 21foo 22fizz 23fizz
池(pool)是一个为处理数量变化并且需要限制并发的greenlet而设计的结构 不使用Pool
1from gevent import monkey 2 3monkey.patch_all() 4import gevent 5import requests 6 7urls = ["https://www.python.org/", "https://www.yahoo.com/", "https://github.com/"] 8 9 10def get(url): 11 print(requests.get(url).url) 12 13 14ts = [gevent.spawn(get, url) for url in urls] 15 16gevent.joinall(ts) 17
使用Pool
1from gevent import monkey 2 3monkey.patch_all() 4import requests 5from gevent.pool import Pool 6 7urls = ["https://www.python.org/", "https://www.yahoo.com/", "https://github.com/"] 8 9 10def get(url): 11 print(requests.get(url).url) 12 13 14p = Pool(3) 15p.map(get, urls) 16
构造一个socket池的类,在各个socket上轮询
1from gevent.pool import Pool 2class SocketPool(object): 3 def __init__(self): 4 self.pool = Pool(10) 5 self.pool.start() 6 def listen(self, socket): 7 while True: 8 socket.recv() 9 def add_handler(self, socket): 10 if self.pool.full(): 11 raise Exception("At maximum pool size") 12 else: 13 self.pool.spawn(self.listen, socket) 14 def shutdown(self): 15 self.pool.kill()
锁和信号量
信号量是一个允许greenlet相互合作,限制并发访问或运行的低层次的同步原语。 信号量有两个方法,acquire和release。在信号量是否已经被 acquire或release,和拥有资源的数量之间不同,被称为此信号量的范围 (the bound of the semaphore)。如果一个信号量的范围已经降低到0,它会 阻塞acquire操作直到另一个已经获得信号量的greenlet作出释放
1from gevent import sleep 2from gevent.pool import Pool 3from gevent.coros import BoundedSemaphore 4sem = BoundedSemaphore(2) 5def worker1(n): 6 sem.acquire() 7 print('Worker %i acquired semaphore' % n) 8 sleep(0) 9 sem.release() 10 print('Worker %i released semaphore' % n) 11def worker2(n): 12 with sem: 13 print('Worker %i acquired semaphore' % n) 14 sleep(0) 15 print('Worker %i released semaphore' % n) 16pool = Pool() 17pool.map(worker1, xrange(0,2)) 18 19运行结果: 20 21Worker 0 acquired semaphore 22Worker 1 acquired semaphore 23Worker 0 released semaphore 24Worker 1 released semaphore
锁(lock)是范围为1的信号量。它向单个greenlet提供了互斥访问。 信号量和锁常被用来保证资源只在程序上下文被单次使用
线程局部变量
Gevent允许程序员指定局部于greenlet上下文的数据。 在内部,它被实现为以greenlet的getcurrent()为键, 在一个私有命名空间寻址的全局查找
1import gevent 2from gevent.local import local 3stash = local() 4def f1(): 5 stash.x = 1 6 print(stash.x) 7def f2(): 8 stash.y = 2 9 print(stash.y) 10 try: 11 stash.x 12 except AttributeError: 13 print("x is not local to f2") 14g1 = gevent.spawn(f1) 15g2 = gevent.spawn(f2) 16gevent.joinall([g1, g2]) 17 18运行结果: 19 201 212 22x is not local to f2
很多集成了gevent的web框架将HTTP会话对象以线程局部变量的方式存储在gevent内。 例如使用Werkzeug实用库和它的proxy对象,我们可以创建Flask风格的请求对象
1from gevent.local import local 2from werkzeug.local import LocalProxy 3from werkzeug.wrappers import Request 4from contextlib import contextmanager 5from gevent.wsgi import WSGIServer 6_requests = local() 7request = LocalProxy(lambda: _requests.request) 8@contextmanager 9def sessionmanager(environ): 10 _requests.request = Request(environ) 11 yield 12 _requests.request = None 13def logic(): 14 return "Hello " + request.remote_addr 15def application(environ, start_response): 16 status = '200 OK' 17 with sessionmanager(environ): 18 body = logic() 19 headers = [ 20 ('Content-Type', 'text/html') 21 ] 22 start_response(status, headers) 23 return [body] 24 WSGIServer(('', 8000), application).serve_forever()
子进程
从gevent 1.0起,支持gevent.subprocess,支持协作式的等待子进程
1import gevent 2from gevent.subprocess import Popen, PIPE 3def cron(): 4 while True: 5 print("cron") 6 gevent.sleep(0.2) 7g = gevent.spawn(cron) 8sub = Popen(['sleep 1; uname'], stdout=PIPE, shell=True) 9out, err = sub.communicate() 10g.kill() 11print(out.rstrip()) 12 13运行结果: 14 15cron 16cron 17cron 18cron 19cron 20Linux
很多人也想将gevent和multiprocessing一起使用。最明显的挑战之一 就是multiprocessing提供
1import gevent 2from multiprocessing import Process, Pipe 3from gevent.socket import wait_read, wait_write 4# To Process 5a, b = Pipe() 6# From Process 7c, d = Pipe() 8def relay(): 9 for i in xrange(5): 10 msg = b.recv() 11 c.send(msg + " in " + str(i)) 12def put_msg(): 13 for i in xrange(5): 14 wait_write(a.fileno()) 15 a.send('hi') 16def get_msg(): 17 for i in xrange(5): 18 wait_read(d.fileno()) 19 print(d.recv()) 20if __name__ == '__main__': 21 proc = Process(target=relay) 22 proc.start() 23 g1 = gevent.spawn(get_msg) 24 g2 = gevent.spawn(put_msg) 25 gevent.joinall([g1, g2], timeout=1) 26 27执行结果: 28 29hi in 0 30hi in 1 31hi in 2 32hi in 3 33hi in 4
然而要注意,组合multiprocessing和gevent必定带来 依赖于操作系统(os-dependent)的缺陷,其中有:
在兼容POSIX的系统创建子进程(forking)之后, 在子进程的gevent的状态是不适定的(ill-posed)。一个副作用就是, multiprocessing.Process创建之前的greenlet创建动作,会在父进程和子进程两方都运行。
上例的put_msg()中的a.send()可能依然非协作式地阻塞调用的线程:一个 ready-to-write事件只保证写了一个byte。在尝试写完成之前底下的buffer可能是满的。
上面表示的基于wait_write()/wait_read()的方法在Windows上不工作 (IOError: 3 is not a socket (files are not supported)),因为Windows不能监视 pipe事件。
Python包gipc以大体上透明的方式在 兼容POSIX系统和Windows上克服了这些挑战。它提供了gevent感知的基于 multiprocessing.Process的子进程和gevent基于pipe的协作式进程间通信
Actors
actor模型是一个由于Erlang变得普及的更高层的并发模型。 简单的说它的主要思想就是许多个独立的Actor,每个Actor有一个可以从 其它Actor接收消息的收件箱。Actor内部的主循环遍历它收到的消息,并根据它期望的行为来采取行动。
Gevent没有原生的Actor类型,但在一个子类化的Greenlet内使用队列, 我们可以定义一个非常简单的
1import gevent 2from gevent.queue import Queue 3class Actor(gevent.Greenlet): 4 def __init__(self): 5 self.inbox = Queue() 6 Greenlet.__init__(self) 7 def receive(self, message): 8 """ 9 Define in your subclass. 10 """ 11 raise NotImplemented() 12 def _run(self): 13 self.running = True 14 while self.running: 15 message = self.inbox.get() 16 self.receive(message)
下面是一个使用的例子:
1import gevent 2from gevent.queue import Queue 3from gevent import Greenlet 4class Pinger(Actor): 5 def receive(self, message): 6 print(message) 7 pong.inbox.put('ping') 8 gevent.sleep(0) 9class Ponger(Actor): 10 def receive(self, message): 11 print(message) 12 ping.inbox.put('pong') 13 gevent.sleep(0) 14ping = Pinger() 15pong = Ponger() 16ping.start() 17pong.start() 18ping.inbox.put('start') 19gevent.joinall([ping, pong])
实际应用
简单server
1# On Unix: Access with ``$ nc 127.0.0.1 5000`` 2# On Window: Access with ``$ telnet 127.0.0.1 5000`` 3from gevent.server import StreamServer 4def handle(socket, address): 5 socket.send("Hello from a telnet!\n") 6 for i in range(5): 7 socket.send(str(i) + '\n') 8 socket.close() 9server = StreamServer(('127.0.0.1', 5000), handle) 10server.serve_forever()
WSGI Servers And Websockets
Gevent为HTTP内容服务提供了两种WSGI server。从今以后就称为 wsgi和pywsgi
- gevent.wsgi.WSGIServer
- gevent.pywsgi.WSGIServer
glb中使用
1import click 2from flask import Flask 3from gevent.pywsgi import WSGIServer 4from geventwebsocket.handler import WebSocketHandler 5import v1 6from .settings import Config 7from .sockethandler import handle_websocket 8def create_app(config=None): 9 app = Flask(__name__, static_folder='static') 10 if config: 11 app.config.update(config) 12 else: 13 app.config.from_object(Config) 14 app.register_blueprint( 15 v1.bp, 16 url_prefix='/v1') 17 return app 18def wsgi_app(environ, start_response): 19 path = environ['PATH_INFO'] 20 if path == '/websocket': 21 handle_websocket(environ['wsgi.websocket']) 22 else: 23 return create_app()(environ, start_response) 24@click.command() 25@click.option('-h', '--host_port', type=(unicode, int), 26 default=('0.0.0.0', 5000), help='Host and port of server.') 27@click.option('-r', '--redis', type=(unicode, int, int), 28 default=('127.0.0.1', 6379, 0), 29 help='Redis url of server.') 30@click.option('-p', '--port_range', type=(int, int), 31 default=(50000, 61000), 32 help='Port range to be assigned.') 33def manage(host_port, redis=None, port_range=None): 34 Config.REDIS_URL = 'redis://%s:%s/%s' % redis 35 Config.PORT_RANGE = port_range 36 http_server = WSGIServer(host_port, 37 wsgi_app, handler_class=WebSocketHandler) 38 print '----GLB Server run at %s:%s-----' % host_port 39 print '----Redis Server run at %s:%s:%s-----' % redis 40 http_server.serve_forever()
3、缺陷
和其他异步I/O框架一样,gevent也有一些缺陷:
- 阻塞(真正的阻塞,在内核级别)在程序中的某个地方停止了所有的东西.这很像C代码中monkey patch没有生效
- 保持CPU处于繁忙状态.greenlet不是抢占式的,这可能导致其他greenlet不会被调度.
- 在greenlet之间存在死锁的可能.
一个gevent回避的缺陷是,你几乎不会碰到一个和异步无关的Python库–它将阻塞你的应用程序,因为纯Python库使用的是monkey patch的stdlib