使用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
解决方案
核心问题分析
BrokenPipeError不触发的原因:TCP协议中,客户端主动断开后,服务器第一次写操作会收到RST包,但如果内核缓冲区还有剩余空间,后续写操作会先将数据写入缓冲区,不会立刻报错,直到缓冲区满或触发TCP超时重传,这个过程可能长达十几分钟。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
相关产品推荐
相关产品推荐

