协程如何工作(How it works?)

原文

概念

进程、线程、协程

一个进程下可以有多个线程, 一个线程又有多个协程, 似乎就是一个从上到下的层次的结构. 但这个并不严谨, 如果能深入协程的允许方式, 就是发现协程与另外两者完全不同, 因为比起由OS调度的进程和线程来说. 协程的调度行为与OS无关.

也就是说, 协程的调度行为完全是在用户空间中进行的, 如果将OS看做一个底层的抽象层, 那么拥有协程的程序只是一个非常普通的程序. 对OS来说, 协程并不真实存在, 一个有拥有协程的线程与其他线程没有本质的区别.

虽然计算机总是在分层, 每一层有对上层提供抽象, 例如进程和线程就是操作系统对用户提供的抽象, 而且线程很多现代编程语言(或者一些古老语言的新标准)都拥有了native的协程支持(例如最常见的语法糖asnc, await), 但这些支持的底层逻辑仍然允许在用户空间, 与操作系统无关.

这种完全在用户空间进行的特性, 是协程与其他两者之间最主要的区别.

协程的抽象

协程常与进程、线程被一并提起, 好像他们是同一种概念的不同层级, 抽象来看确实如此。 我们首先就以这样抽象的角度区审视和感受线程的表现,而不关注他们的细节。

协程与一个函数类似, 不过函数调用总是从一个入口进入, 从一个出口返回, 而协程则不同。 协程虽然类似一个函数, 但在执行过程中可以中断, 被调度区执行其他协程或普通函数, 然后在未来某个时刻调度回来继续执行。

如果把协程作为一个黑箱, 那么在运行时, 它看起来就是一个独立调度的单位,但与进程、线程这种独立调度不同的是, 协程挂起和恢复的位置总是固定的,而且是用户开发组显示指定的。

为了更好的理解这些特性,我们用线程做对比:

这里使用js作为描述语言, 忽略v8引擎让js执行在单线程中.

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15

let sum = 0;
function thread1(){
for(let _ in [...new Array(10000).keys()]){
// do sth.
sum += 1;
}
}

function thread2(){
for(let _ in [... new Array(10000).keys()]){
// do sth.
sum += 1;
}
}

使用thread1, thread2两个线程并反复执行, 敏感的开发者可以轻易发现存在的问题—对sum的访问在线程的并发环境下不是原子的, 这会导致sum的最终值不一定是10000+10000=20000, 因为OS很可能在非原子操作的间隙调度线程, 导致上下文变量在暂停和恢复时拥有错误的值.

但如果使用协程:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15

let sum = 0;
function coroutine1() {
for(let _ in [...new Array(10000).keys()]) {
await sth;
sum += 1;
}
}

function coroutine2() {
for(let _ in [...new Array(10000).keys()]) {
await sth;
sum += 1;
}
}

我们执行这两个协程, 最终sum的值一定会停在20000, 这是因为, 协程只会在一个被显式指定的位置被调度,在上面的例子中, 就是在await sth; 的位置被调度, 这就可以保证每个原子操作都会完整地执行.

这里没有体现协程的执行顺序, 这个将在后面的例子中说明

这就是协程的抽线表现, 看起来让人有些费解, 因为如果能指定协程中被调度的位置, 意味着调度器需要与协程沟通, 知晓可以在哪个位置调度它, 也就是说, 调度器与每个协程都是耦合的.

OS的并发问题就来源于OS不知道进程或线程内部的抽象结构, 因为一个高级语言中的本该是原子操作的代码被编译到机器码后, 就不再是一个原子操作了, 它可能包含多条机器指令, 这就意味着机器码
丢失了高级语言的这些抽象信息, 从而可能在连续的多条机器指令间隙被调度, 造成并发环境下的逻辑混乱.

我们当然可以想象, 如果能与OS沟通, 让OS知道哪个位置的代码应该是原子的, 那么就可以解决并发问题, 而这个解决方案就是基于信号量的进程同步, 也可以说是.

但这又与协程完全不同, OS仅仅是执行锁的原子操作, 而协程将可以指定几个调度点位, 告诉调度器可以在哪个点位进行调度.

就好像把一个函数截成好几段, 每执行一段后调度器就接入, 并进行调度, 将来还会从该位置恢复上下文.

除了可以把函数分段进行调度外, 协程真正的威力主要来源于调度和异步I/O库的配合使用, 我们可以让协程在I/O时运行被调度, 从而充分利用I/O时间来运行其他CPU密集的代码.

这是如何做到的呢?

协程的调度:

现在我们知道了:

  • 协程的调度完全在用户空间中进行
  • 协程只能从被显式声明的可调度的位置被调度

前缀说明了线程的普通的性质, 协程的调度器只是某个线程中的某一段无聊的逻辑, 并没有什么特别的, 而后者则是下面我们将要了解的问题.

如何实现这个调度器, 才能让一个独立的调度单位于调度器沟通可调度的位置, 并允许一个函数在调度位置中途离开, 在将来恢复?

实际上目前大多数现代编程语言都提供了实现协程调度器的一个基本特性, 那就是—生成器(generator).

在Python中, 这被成为生成器迭代器, 由生成器迭代器函数返回:

任何一个拥有yield表达式的普通函数都将被解释器处理为一个生成器迭代器函数, 该函数执行后将返回一个生成器迭代器:

1
2
3
4
5
6
7
8
9
10

def gen():
yield 1
# do sth
yield 2
# do sth

coro = gen()
coro.send(None) # 1
coro.send(None) # 2

每次执行send方法,将会把一个值传入这个迭代器内部, 然后一直执行到下一个yield语句的位置, 将该语句后的表达式返回值从刚才的send调用处返回。

这是不是很像协程呢? 每一个yield语句就相当于协程的一次分段,在每个分段处,生成器迭代器总是将控制权交还初期, 我们完全可以实现一个调度器, 这个调度器调度很多生成器迭代器,每个生成器迭代器抽象上来说就是一个完全的协程。

协程的优势

协程的应用场景几乎只有一个, 那就是IO密集的程序, 例如典型的高IO场景—Web后端服务

  • 开销小
    • 协程调度的性能开销比线程调度小很多, 协程切换迅速
    • 协程调度的内存开销也比线程调度小很多
  • 一般情况下,不需要考虑在并发环境下对资源原子性访问的问题, 节省了锁的开销
  • 逻辑更简洁
    • 相比线程来说, 协程让从上下文推导代码逻辑变得更加简单

Threads make local reasoning difficult, and local reasoning is perhaps the most importantthing in software develpoment.

线程使得局部推理变得困难, 而军部推理可能是软件开发中最重要的事情。

Node.js就是一个基于单线程非阻塞I/O模型开发的JavaScript Runtime, 只要实现得当, 即便是单线程模型, 也可以承载相当巨大的并发量。

举一个简单的例子, 例如一个典型的爬虫程序:

1
2
3
4
import requests
for i in range(100):
resp = requests.get(url="http:gaolihai.cool/doc/README.md")
print("\n".join([linebytes.decode() for linebytes in resp.iter_lines()]))

假设每次从请求到返回需要100s, 假设其中发送请求花费了5ms, 等待服务端响应, 即等待IO结束花费了95ms, 那么100次请求将花费10s, 而这10s中只有0.5s即发送请求的过程是必须占用CPU时间的, 剩下9.5s只是无意义的等待。 这种模型被称为阻塞式I/O或同步I/O。

而如果使用协程, 我们可以这样实现:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
from asynclib import asynchttp, Future

@asyncfun
def http():
responseData = yield from Future(
lambda resolve:
asynchttp.get(
url='http://gaolihai.cool/doc/README.md',
callback=lambda response: resolve(response)
)
)
print(responseData)

for i in range(100):
http()

asynchttp.get是一个异步请求api, 即它只负责发送请求, 所以还额外接受一个回调函数作为参数, 该函数将在IO结束时由操作系统通知进程, 然后进行回调。

那么这段程序只需要约0.5s就能发送完所有请求, 这可以称为并发请求, 等待服务端响应可能需要200ms,最终在不到一秒内就完成了100次请求。这种模型被称为非阻塞式I/O或异步I/O。

相比阻塞式I/O, 非阻塞式I/O在这种场景下显然具有很大优势。 不过如果我们使用线程, 也能得到差不多的性能表现, 只不过线程的开销相比协程要大很多,当并发数量过多时可能会在线程上下文切换中浪费太多性能, 而过多的线程实例甚至会挤爆栈空间。

所以, 协程本身并不会让代码执行很快, 但如果配合异步I/O库, 在高I/O场景下将具有非常大的性能优势。

另外, 对于CPU密集的成功徐来说, 更建议使用多进程, 在多个核心上进行并行计算。

实现及原理

协程库实现

协程调度的实现方式由很多, 这里介绍一种使用事件循环、是按队列这种调度方式的协程实现。

事件循环和事件队列

一个典型的Python程序:

1
2
3
import requests
resp = requests.get(url='http://gaolihai.cool/doc/README.md')
print('\n'.join([linebytes.decode() for linebytes in resp.iter_lines()]))

一个典型的使用协程的程序:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
from asynclib.core import Future, asyncRun, loop
from asynclib.asynchttp import get as asyncget

def http():
responseData = yield from Future(
lambda resolve:
asyncget(
url='http://gaolihai.cool/doc/README.md',
callback=lambda response: resolve(response)
)
)
print(responseData.decode('utf-8'))

asyncRun(http())
loop()

协程调度的典型流程是这样的:

现在有两个实体:

  • 协程调度器: 事件循环是调度器的实现方式
  • 事件队列: 存放所有待执行的事件

asyncRun调用可以将一个协程压入事件队列中, loop是进入事件循环(也可称为调度器)的入口, loop调用将会把线程的控制权交给协程调度器, 该调度器将会在未来不断地从事件队列拉取协程或普通函数(可以称事件), 然后执行和调度它们.

在调度和执行的过程中, 这些是按还可能产生更多的事件, 于是就会源源不断地执行下去.

事件队列使用一个阻塞队列实现, 对外提供一个单例eventQueue:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
from queue import Queue


class __EventQueue:
def __init__(self) -> None:
self.__eventQueue = Queue()

def pushCallback(self, fn):
self.__eventQueue.put(fn, block=True)

def getCallback(self):
return self.__eventQueue.get(block=True)

eventQueue = __EventQueue()

一个最简单的事件循环实现:

当队列中没有元素时, eventQueue.getCallback()调用将会阻塞在内部的self.__eventQueue.get(block=True)处, 直到异步I/O库某个接口执行完毕后, 向该队列中压入一个事件, 此时阻塞在此处的事件循环线程将被唤醒, 然后执行事件.

1
2
3
4
def loop():
while True:
cbk = eventQueue.getCallback()
cbk()

asyncRun的职责是向事件队列中压入一个事件(或者说任务):

1
2
3
def asyncRun(gen):
eventQueue.pushCallback(gen)

到此为止, 一个基本的框架就搭好了, 下面我们来看, 如何进行协程的调度.

生成器迭代器的自动执行

我们刚才提到, 基于生成器迭代器(下文称为协程)来实现协程, 现在就需要一个协程的自动执行器.

在事件循环的实现中, 如果拿到一个协程, 肯定不能简单地调用它, 我们需要一个执行器, 进行这样一个非常关键的流程, 这个流程会反复在协程和调度器之间切换控制权, 从而自动地异步地将协程执行完毕:

协程的自动执行器将实现一个过程, 该过程通过接受一个协程参数, 并调用send()将控制权i奥给协程(相当于调度协程).

在协程执行异步api时, 再从yield处将控制权交给调度器, 转而执行事件队列中下一个任务.

当I/O结束时, 由操作系统通知异步api, 然后api内部将会向事件队列压入一个回调函数, 当函数被事件循环取出和调用后, 再次转移控制权给协程, 持续驱动协程的异步执行, 就这样一直来回递交控制权, 知道协程完全执行完毕.

修改事件循环

首先修改事件循环的实现, 其中__GeneratorExecutor是协程的执行器类

1
2
3
4
5
6
7
8
9
10
11
12
from inspect import isgenerator, isgeneratorfunction
def loop():
while True:
cbk = eventQueue.getCallback()
if isgenerator(cbk):
self.__GeneratorExecutor(cbk)
elif isgeneratorfunction(cbk):
self.__GeneratorExecutor(cbk())
elif callable(cbk):
cbk()
else:
raise TypeError('...')

执行器和Future

在执行器的实现中, 执行器和Future是密不可分的, 因为要驱动协程反复在执行器和协程之间转移空hi全, 必须有一个规定的接口, 这恶鬼接口就是Fture, 我们规定协程在yield后必须跟一个Future对象.

1
class Future:

这个Future将被这样设计:

它于JavaScript的Promise类似,

  • Future构造函数接收一个任务参数, 该参数是一个具有resolve参数的函数, 在函数中调用异步任务, 并在异步任务的回调中调用resolve, 并传入异步任务的结果来通知Future对象.

    1
    2
    3
    4
    def __init__(self, task=None):
    self.callbacks: List[Callable[[Any], Any]] = []
    self.value: Any = None
    task(self.resolve)
  • 允许在Future对象上添加回调函数, 该回调函数将在异步任务结束(即resolve被调用后)由Future对象进行回调. 由于reslve被定义为异步任务结束, 用来通知Future回调, 于是该函数具有一个参数, 就是异步函数完成的返回值.

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    def addCallback(self, cbk: Callable[[Any], Any]):
    if self.statee == 'resolved':
    cbk(self)
    self.callbacks.append(cbk)
    return self

    def resolve(self, value: Any = None):
    self.value = value
    for cbk in self.callbacks:
    cbk(self)
    return self

下main来看执行器:

执行器的实现非常简单, 它的构造函数接收一个协程, 然后通过__next方法直接调度协程.

直到协程返回一个Future后, 调度器拿到控制权, 并在这个返回的Future删注册__next方法作为回调函数, 这样就可以在该Future包装的异步任务执行完毕后自动回调执行器, 交换控制权, 并继续下一步的注册回调函数, 就这样反复交换控制权, 驱动协程执行直到结束.

1
2
3
4
5
6
7
8
9
10
11
12

class __GeneratorExecutor:
def __init__(self, coroutine):
self.coroutine = coroutine
self.__next(Future())

def __next(self, future: Future):
try:
nextFuture = self.coroutine.send(future.value)
except StopIteration:
return
nextFuture.addCallback(self.__next)

此外, 每次执行器将控制权转移给协程时, 将会把上一个Future对象的异步结果值传入协程, 在协程内部来看, 就好像是异步任务在I/O时阻塞, I/O结束后从同步上下文中返回了异步任务的结果.

到此位置, 一个最小功能的协程就实现好了, 现在我们拥有:

  • 事件循环
  • 事件队列
  • Future和协程执行器

下面我们简单实现一个异步HTTP GET请求来测试协程的行为.

异步api举例

下面我们要基于非阻塞socket实现一个简单的HTTP GET请求接口.

非阻塞socket拥有三个异步过程:

  • 第一阶段是等待连接, 即等待socket文件可写
  • 发送请求后, 第二阶段是等待服务器响应, 即等待socket文件可读

由于是非阻塞I/O, 每一步都需要通过注册回调函数来处理, 且每次socket文件状态改变后由OS通知, 这一点通过selector(I/O多路复用)实现

该接口应该接受两个参数:

  • 要访问的url
  • 服务端响应后的回调函数callback

实例化socket对象, 并设定为非阻塞模式

实例化selector对象, 并监听socket文件的可写事件, 此时还需要实现连接成功后的回调connected

发送连接请求, 连接成功后connected将会首先发送数据, 然后组测刻度事件, 等待服务端响应, 此时需要实现responded

repsonded将会连续接受数据知道遇到数据结束标志, 到此一次从请求到响应的完整的HTTP请求就i结束了, 现在向事件队列压入回调函数, 通知调用者异步任务结束。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
def get(*, url, callback):
urlObj = urllib.parse.urlparse(url)
selector = DefaultSelector()
sock = socket.socket()
sock.setblocking(False)

def connected():
selector.unregister(sock.fileno())
selector.register(sock.fileno(), EVENT_READ, responded)
sock.send(
f"""GET {urlObj.path if urlObj.path != '' else '/'}{'?' if urlObj.query != '' else '' + urlObj.query} HTTP/1.0\r\n\r\n"""
.encode('ascii')
)

responseData = bytes()

def responded():
nonlocal responseData
chunk = sock.recv(4096)
if chunk:
responseData += chunk
else:
selector.unregister(sock.fileno())
eventQueue.pushCallback(lambda: callback(responseData))
nonlocal __stop
__stop = True
__stop = False

def loop():
while True:
events = selector.select()
for event_key, event_mask in events:
cbk = event_key.data
cbk()
if __stop:
break

selector.register(sock.fileno(), EVENT_WRITE, connected)
try:
sock.connect(
(urlObj.hostname, urlObj.port if urlObj.port != None else 80)
)
except BlockingIOError:
pass
Thread(target=loop).start()

与之前的连连起来讲, 这个回调函数在将来被事件循环调用, 如果它在一个协程中使用, 将会resolve外层Future。

随后Future依次调用注册在它身上的回调函数, 其中一个ius执行器的__next方法, 从而将控制权转移到执行器中。

再往后, 执行器将会把异步任务的结果传入协程继续执行, 从而让协程中的该异步任务带着执行器传入的任务结果从yield恢复执行。