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

使用Python http.server实现流端点时,如何可靠检测客户端断开?

如何可靠检测HTTP长连接客户端是否在线?

我们用Python的http.server实现了一个简单Web服务器,其中一个端点提供事件流推送功能,新事件产生时实时推送给客户端。目前数据推送正常,但客户端断开后服务器无法检测,会持续向无效连接发送数据。

最初的代码期望调用self.wfile.write()时触发BrokenPipeError来检测连接断开,但实际该错误不会触发,服务器甚至会持续推送20分钟以上无报错。之后尝试设置socket超时并使用select()检测socket可写性的方案,依然无效:客户端断开数分钟后既无超时错误,select()仍返回socket可写。

现求问:如何正确且可靠地检测客户端是否仍在线监听?


初始代码

import json
import queue
import http.server

from common.log import log

logger = log.get_logger()

class RemoteHTTPHandler(http.server.BaseHTTPRequestHandler):

    ...

    def __stream_events(self, start_after_event_id: int) -> None:
        """Stream events starting after the given ID and continuing as new events become available"""

        # 获取事件队列,包含指定起始点后的所有已有事件,新事件产生时会自动加入队列
        logger.info(f"Streaming events from ID {start_after_event_id}")
        with self._events_stream_manager.stream_events(start_after_event_id) as events_queue:
            self.send_response(200)
            self.send_header("Content-type", "application/yaml; charset=utf-8")
            self.send_header("Connection", "close")
            self.end_headers()

            # 服务器关闭时终止所有流连接
            while not self._stop_streams_event.is_set():
                try:
                    # 获取下一个事件;队列为空时最多阻塞1秒
                    try:
                        data = events_queue.get(timeout=1)
                    except queue.Empty:
                        # 发送空行防止HTTP连接超时
                        self.wfile.write(b"\n")
                        continue

                    # 发送编码后的事件及分隔线
                    self.wfile.write(json.dumps(data, indent=4).encode('utf-8') + b"\n\n---\n\n")
                except BrokenPipeError as ex:
                    # TODO: 无法可靠检测连接断开
                    # 管道破裂表示连接已丢失,可能是客户端主动关闭或网络错误
                    logger.info(f"Connection closed: {type(ex).__name__}: {ex}", exc_info=True)
                    return


def serve():
    http.server.ThreadingHTTPServer(("", 8090), RemoteHTTPHandler).serve_forever()

尝试优化后的代码

import http.server
import json
import queue
import select
import socket

from common.log import log

logger = log.get_logger()


class RemoteHTTPHandler(http.server.BaseHTTPRequestHandler):

    ...

    def __stream_events(self, start_after_event_id: int) -> None:
        """Stream events starting after the given ID and continuing as new events become available"""

        # 获取事件队列,包含指定起始点后的所有已有事件,新事件产生时会自动加入队列
        logger.info(f"Streaming events from ID {start_after_event_id}")
        with self._events_stream_manager.stream_events(start_after_event_id) as events_queue:
            # 发送响应头
            self.send_response(200)
            self.send_header("Content-type", "application/yaml; charset=utf-8")
            self.send_header("Connection", "close")
            self.end_headers()

            # 为底层socket设置超时
            self.connection.settimeout(2)

            # 服务器关闭时终止所有流连接
            while not self._stop_streams_event.is_set():
                try:
                    # 获取下一个事件;队列为空时最多阻塞1秒
                    message: bytes
                    try:
                        data = events_queue.get(timeout=1)
                    except queue.Empty:
                        # 发送空行防止HTTP连接超时
                        logger.debug(f"Sending blank line")  # FIXME: 测试用
                        message = b"\n"
                    else:
                        # 发送编码后的事件及分隔线
                        logger.debug(f"Sending event")  # FIXME: 测试用
                        message = json.dumps(data, indent=4).encode('utf-8') + b"\n\n---\n\n"

                    # 用select检测客户端是否仍连接
                    write_ready_fds: list[socket.socket]
                    logger.debug(f"Connection check")  # FIXME: 测试用
                    __, write_ready_fds, __ = select.select([], [self.connection], [], 3)
                    if not write_ready_fds:
                        logger.info(f"Connection closed")
                        return

                    logger.debug(f"Still connected; sending")  # FIXME: 测试用
                    self.wfile.write(message)
                except (BrokenPipeError, TimeoutError) as ex:
                    # 管道破裂或socket超时表示连接已丢失,可能是客户端主动关闭或网络错误
                    logger.info(f"Connection closed: {type(ex).__name__}: {ex}", exc_info=True)
                    return
                except BaseException as ex:
                    # 意外错误
                    logger.exception(f"Error sending events: {type(ex).__name__}: {ex}")
                    raise

解决方案

核心问题分析

  1. BrokenPipeError不触发的原因:TCP协议中,客户端主动断开后,服务器第一次写操作会收到RST包,但如果内核缓冲区还有剩余空间,后续写操作会先将数据写入缓冲区,不会立刻报错,直到缓冲区满或触发TCP超时重传,这个过程可能长达十几分钟。
  2. select()检测无效的原因:select()返回可写仅表示socket内核缓冲区有空闲空间,不代表客户端仍在接收数据。客户端断开后只要缓冲区未填满,select()依然会返回可写状态。

可靠检测方案

方案1:直接操作底层socket并检查发送结果

绕过BufferedWriter的缓冲区,直接调用系统级send函数发送数据,同时检查发送字节数是否完整:

import os

# ... 其他代码 ...

# 替换原有的self.wfile.write(message)逻辑
try:
    # 直接通过socket文件描述符发送数据
    sent_bytes = os.send(self.connection.fileno(), message)
    if sent_bytes != len(message):
        # 未发送完整,说明连接已断开
        logger.info("Connection closed: partial data send failed")
        return
except BrokenPipeError:
    logger.info("Connection closed: BrokenPipeError detected")
    return

方案2:启用TCP Keepalive(推荐)

利用系统级TCP心跳机制主动检测连接状态,这是最可靠的方案,无需修改应用层逻辑:

def __stream_events(self, start_after_event_id: int) -> None:
    # ... 发送响应头之后的位置 ...
    # 开启TCP Keepalive
    self.connection.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
    # 设置空闲10秒后开始发送心跳包
    self.connection.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPIDLE, 10)
    # 心跳包发送间隔3秒
    self.connection.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPINTVL, 3)
    # 连续3次心跳无响应则判定连接断开
    self.connection.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPCNT, 3)

    # ... 后续循环逻辑保持不变 ...

启用后,系统会自动定期发送心跳包检测连接,若客户端断开,内核会触发错误,此时写操作会立刻抛出BrokenPipeError。

方案3:应用层心跳探测(需客户端配合)

在队列空时,发送一个约定的心跳包(比如X-Heartbeat标识的HTTP分块),要求客户端响应。但该方案需要客户端配合实现,通用性不如TCP Keepalive。

推荐方案

优先选择方案2(TCP Keepalive),这是系统级的原生检测机制,可靠性最高,无需额外应用层逻辑。同时可以搭配方案1,直接操作底层socket发送数据,避免缓冲区延迟报错。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 14:54:53