如何同时等待Socket Pollin事件与Queue消息触发?
问题描述
我有一个用于监听Socket传入消息的select.poll()对象,还有一个存储待发送消息的queue.Queue()对象。两者各自支持带超时的无CPU消耗等待,但我需要同时等待这两个对象,任一触发就恢复线程执行,有没有对应的实现机制?
我试过在循环中交替用极短超时等待两者,但本质是忙等待,CPU占用太高;增加超时能降低CPU占用,但又会大幅提升延迟。
因为select.poll()可以监听任意UNIX文件描述符,我的问题也可以换成:有没有能暴露可轮询文件描述符的queue.Queue()或threading.Event()替代方案?
当前实现
import socket import select import queue import threading socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) socket.connect(('localhost', 4200)) incoming = select.poll() incoming.register(socket, select.POLLIN) outgoing = queue.Queue() while True: if incoming.poll(timeout=0.001): read(socket) try: msg = outgoing.get(block=False, timeout=0.001) send(socket, msg) except queue.Empty: pass
期望实现
import socket import queue import threading socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) socket.connect(('localhost', 4200)) outgoing = queue.Queue() poller = HYBRID_POLLER() # TODO poller.register(socket) poller.register(outgoing) while True: available = poller.poll(timeout=1) if socket in available: read(socket) if outgoing in available: msg = outgoing.get(block=False, timeout=0.001) send(socket, msg)
解决方案:用管道将队列事件转为可轮询的文件描述符
在UNIX-like系统下,可以利用**管道(pipe)**实现需求:当队列有消息时,往管道写端写入一个字节;让select.poll()监听管道读端,这样队列有消息时管道读端会触发POLLIN事件,和Socket事件一起被poll捕获,实现无CPU消耗的同时等待。
实现代码
import socket import select import queue import os import threading def read(sock): data = sock.recv(1024) print(f"Received: {data.decode()}") def send(sock, msg): sock.send(msg.encode()) print(f"Sent: {msg}") # 创建管道,用于将队列事件转为可poll的文件描述符 pipe_r, pipe_w = os.pipe() # 封装队列,添加消息时自动触发管道事件 class PollableQueue(queue.Queue): def put(self, item, block=True, timeout=None): super().put(item, block, timeout) # 往管道写一个字节,触发POLLIN事件 os.write(pipe_w, b'x') socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) socket.connect(('localhost', 4200)) # 初始化poller,注册Socket和管道读端 poller = select.poll() poller.register(socket, select.POLLIN) poller.register(pipe_r, select.POLLIN) outgoing = PollableQueue() # 模拟线程往队列添加消息 def producer(): import time while True: outgoing.put(f"Hello at {time.time()}") time.sleep(2) threading.Thread(target=producer, daemon=True).start() while True: events = poller.poll(timeout=-1) # 无限等待直到有事件触发 for fd, event in events: if fd == socket.fileno(): read(socket) elif fd == pipe_r: # 读取管道中所有标记,避免缓存溢出导致重复触发 os.read(pipe_r, 1024) # 批量处理队列中的所有消息 while True: try: msg = outgoing.get(block=False) send(socket, msg) except queue.Empty: break
原理说明
- 管道的核心作用:管道读端是可被
select.poll()监听的文件描述符,当管道有数据可读时,会触发POLLIN事件。 - 队列封装逻辑:自定义
PollableQueue继承自queue.Queue,重写put()方法,每次添加消息后往管道写端写入一个字节,主动触发管道读端的事件。 - 事件处理流程:poll循环中检测到管道事件时,先清空管道缓存,再批量处理队列中的所有消息,避免遗漏或重复触发。
注意事项
- 该方案仅适用于UNIX-like系统(Linux、macOS等),Windows系统不支持
os.pipe()返回的文件描述符被select.poll()监听。 - 处理管道事件时必须一次性读取所有数据,否则残留字节会导致后续重复触发事件。
- 若需支持Windows,可改用
win32pipe创建命名管道,或直接切换到asyncio框架(原生支持同时等待Socket和队列/事件)。
内容的提问来源于stack exchange,提问作者danijar
相关产品推荐
相关产品推荐

