David Beazley协程斐波那契服务器中future_monitor()与socket pair的必要性
future_monitor()的作用解析 在PyCon 2015的演讲中,展示了一段基于协程的斐波那契服务器代码,必须添加future_monitor()才能使服务器正常运行(否则会陷入卡顿)。核心代码如下:
from socket import * from collections import deque from concurrent.futures import ProcessPoolExecutor as Pool from select import select pool = Pool(4) def fib(n): if n <= 2: return 1 else: return fib(n-1) + fib(n-2) def fib_server(address): sock = socket(AF_INET, SOCK_STREAM) sock.setsockopt(SOL_SOCKET, SO_REUSEADDR, 1) sock.bind(address) sock.listen(5) while True: yield 'recv', sock conn, addr = sock.accept() # blocking print("connection", addr) tasks.append(fib_handler(conn)) def fib_handler(conn): while True: yield 'recv', conn req = conn.recv(100) # blocking if not req: break n = int(req) future = pool.submit(fib, n) yield 'future', future result = future.result() # blocking resp = str(result).encode('ascii') + b'\n' yield 'send', conn conn.send(resp) # blocking print('closed') tasks = deque() recv_wait = {} send_wait = {} future_wait = {} future_notify, future_event = socketpair() def future_done(future): tasks.append(future_wait.pop(future)) future_notify.send(b'x') def future_monitor(): while True: yield 'recv', future_event future_event.recv(100) tasks.append(future_monitor()) def run(): while any([tasks, recv_wait, send_wait]): while not tasks: # no active task to run wait for IO can_recv, can_send, _ = select(recv_wait, send_wait, []) for s in can_recv: tasks.append(recv_wait.pop(s)) for s in can_send: tasks.append(send_wait.pop(s)) task = tasks.popleft() try: why, what = next(task) if why == 'recv': recv_wait[what] = task elif why == 'send': send_wait[what] = task elif why == 'future': future_wait[what] = task what.add_done_callback(future_done) else: raise RuntimeError("We don't know what to do with :", why) except StopIteration: print('task done') if __name__ == "__main__": tasks.append(fib_server(('localhost', 25000))) run()
技术疑问
- 为何需要
future_monitor()? - 仅使用
future_done()回调(不写入socket pair)为何无法满足需求? - 为何必须使用socket pair?测试发现移除socket pair后服务器会失去响应。
核心解答
1. 单独用future_done()回调为什么不行?
这个服务器的事件循环依赖select()系统调用阻塞等待IO事件,只有当select()检测到可读/可写的socket时,才会唤醒循环处理任务队列。
当ProcessPoolExecutor的子进程完成任务后,future_done()是由concurrent.futures的内部线程触发执行的——此时主进程的事件循环可能正卡在select()调用上,完全没在检查tasks队列。哪怕你在回调里把任务放回tasks,事件循环也感知不到,只会一直卡在select(),服务器自然卡顿。
2. future_monitor()的作用是什么?
future_monitor()是一个专门的协程,它的唯一作用就是把线程回调的通知转换成select()能识别的IO事件。它一直监听future_event这个socket的可读状态,当future_done()往配对的future_notify写数据时,future_event会变成可读,select()会立刻检测到这个变化,唤醒事件循环去处理future_monitor()协程,此时事件循环才会重新检查tasks队列,发现刚放回的任务并继续执行。
3. 为什么必须用socket pair?
select()只能监听**文件描述符(包括socket)**的IO事件,没法直接感知线程回调或future的完成状态。socket pair是一对互相连通的socket,往其中一个写数据,另一个就会触发可读事件,刚好能被select()捕获。
这种方式相当于给线程回调和主事件循环之间搭了个通信桥梁,把异步任务完成的信号,转换成了事件循环能识别的IO事件。如果移除socket pair,future_done()放回任务后,事件循环永远卡在select(),根本不知道有新任务,服务器必然失去响应。
内容的提问来源于stack exchange,提问作者Ben noop

