You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 14:22:04