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

如何同时等待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

原理说明

  1. 管道的核心作用:管道读端是可被select.poll()监听的文件描述符,当管道有数据可读时,会触发POLLIN事件。
  2. 队列封装逻辑:自定义PollableQueue继承自queue.Queue,重写put()方法,每次添加消息后往管道写端写入一个字节,主动触发管道读端的事件。
  3. 事件处理流程:poll循环中检测到管道事件时,先清空管道缓存,再批量处理队列中的所有消息,避免遗漏或重复触发。

注意事项

  • 该方案仅适用于UNIX-like系统(Linux、macOS等),Windows系统不支持os.pipe()返回的文件描述符被select.poll()监听。
  • 处理管道事件时必须一次性读取所有数据,否则残留字节会导致后续重复触发事件。
  • 若需支持Windows,可改用win32pipe创建命名管道,或直接切换到asyncio框架(原生支持同时等待Socket和队列/事件)。

内容的提问来源于stack exchange,提问作者danijar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 04:57:03