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

如何让Asyncio Streams服务器检测客户端网络断开?

问题

我正在使用Asyncio Streams重写旧服务器,其中一直未能解决断连检测问题。我编写了简单的服务器示例:

import asyncio

async def handle_client(reader, writer):
    address = writer.get_extra_info('peername')
    print(f"Connection: {address}")

    try:
        while True:
            data = await reader.read(1024)
            if not data:
                print(f"Disconnected: {address}.")
                break
            print(f"{address} Send: {data.decode()}")

            writer.write('👍\n'.encode())
            await writer.drain()
    except Exception as e:
        print(f"Something bad happened: {e}")
    finally:
        print(f"Closed: {address}")
        writer.close()
        await writer.wait_closed()


async def main():
    server = await asyncio.start_server(handle_client, '0.0.0.0', 5678)

    async with server:
        await server.serve_forever()

asyncio.run(main())

使用netcat连接时一切正常,关闭netcat或终端时服务器能检测到断连,但客户端断开网络时,服务器无法感知。我知道这是底层限制,目前采用超时强制断连的方案:

import asyncio

async def handle_client(reader, writer):
    address = writer.get_extra_info('peername')
    print(f"Connection: {address}")

    try:
        while True:
            data = await asyncio.wait_for(reader.read(1024), timeout=10.0)
            if not data:
                print(f"Disconnected: {address}.")
                break
            print(f"{address} Send: {data.decode()}")

            writer.write('👍\n'.encode())
            await writer.drain()
    except asyncio.TimeoutError:
        print(f"Timeout: {address}")
    except Exception as e:
        print(f"Something bad happened: {e}")
    finally:
        print(f"Closed: {address}")
        writer.close()
        await writer.wait_closed()

async def main():
    server = await asyncio.start_server(handle_client, '0.0.0.0', 5678)

    async with server:
        await server.serve_forever()

asyncio.run(main())

该方案有效,但需要针对不同设备调整参数,而我的多数客户端因历史原因心跳间隔极长。请问有没有无需等待数分钟超时就能检测断连的更优方法?

解决方案

1. 启用TCP Keepalive机制

TCP原生支持Keepalive功能,可通过操作系统底层配置或直接在Asyncio的socket上设置,无需修改应用层逻辑就能快速检测死连接。

修改handle_client函数,在连接建立后配置socket的Keepalive参数:

import asyncio
import socket

async def handle_client(reader, writer):
    address = writer.get_extra_info('peername')
    # 获取底层socket并开启Keepalive
    sock = writer.get_extra_info('socket')
    sock.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
    # 设置首次探测前的空闲时间(秒)
    sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPIDLE, 5)
    # 设置探测失败的重试次数
    sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPCNT, 3)
    # 设置每次探测的间隔(秒)
    sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPINTVL, 2)
    
    print(f"Connection: {address}")
    try:
        while True:
            data = await reader.read(1024)
            if not data:
                print(f"Disconnected: {address}.")
                break
            print(f"{address} Send: {data.decode()}")
            writer.write('👍\n'.encode())
            await writer.drain()
    except Exception as e:
        print(f"Something bad happened: {e}")
    finally:
        print(f"Closed: {address}")
        writer.close()
        await writer.wait_closed()

当客户端断网后,TCP层会自动发送探测包,多次失败后会触发连接断开,此时reader.read()会返回空或抛出异常,服务器能立刻感知断连。该方法依赖操作系统支持,主流Linux、Windows、macOS均兼容。

2. 分离业务读取与断连检测超时

如果客户端心跳间隔长,但需要快速检测断连,可以用asyncio.wait同时监听业务读取和短周期的断连探测任务,避免单一长超时导致的延迟:

import asyncio
import socket

async def handle_client(reader, writer):
    address = writer.get_extra_info('peername')
    # 先开启TCP Keepalive作为基础保障
    sock = writer.get_extra_info('socket')
    sock.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
    sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPIDLE, 10)
    sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPCNT, 3)
    sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPINTVL, 2)
    
    print(f"Connection: {address}")
    try:
        while True:
            # 同时等待读取数据和30秒超时任务,谁先完成就处理谁
            done, pending = await asyncio.wait(
                [reader.read(1024), asyncio.sleep(30)],
                return_when=asyncio.FIRST_COMPLETED
            )
            # 处理读取到的业务数据
            if reader.read(1024) in done:
                data_task = done.pop()
                data = data_task.result()
                if not data:
                    print(f"Disconnected: {address}.")
                    break
                print(f"{address} Send: {data.decode()}")
                writer.write('👍\n'.encode())
                await writer.drain()
                # 取消未完成的超时任务,进入下一轮循环
                for task in pending:
                    task.cancel()
            else:
                # 超时触发,主动发送空包触发TCP层ACK检测
                try:
                    writer.write(b'')
                    await writer.drain()
                except Exception as e:
                    print(f"Detected dead connection: {address}, error: {e}")
                    break
    except Exception as e:
        print(f"Something bad happened: {e}")
    finally:
        print(f"Closed: {address}")
        writer.close()
        await writer.wait_closed()

这个逻辑的核心是:每隔30秒主动发送一个空包(不影响业务),如果客户端断连,writer.drain()会抛出异常,从而立刻检测到断连。只要客户端有业务数据发送,就会重置超时检测,不会干扰长心跳的客户端。

3. 应用层自适应心跳协商(可选)

如果允许修改客户端逻辑,可在连接建立时让客户端上报自身的心跳间隔,服务器根据每个客户端的参数动态调整超时时间。但如果客户端无法修改,此方法不适用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 17:02:18