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

如何实现可通过外部事件关闭的跨平台Python WebSocket服务器?

跨平台Python WebSocket服务器优雅关闭方案

我来帮你解决这个跨平台关闭WebSocket服务器的问题。首先咱们拆解下你当前代码的核心问题,再给出通用的解决方案——毕竟跨平台兼容(尤其是Windows不支持add_signal_handler)是关键。

当前代码的核心问题

你用loop.call_soon_threadsafe(stop_listening, loop, server)触发关闭,但stop_listening是异步函数,call_soon_threadsafe只能调度同步函数,没法处理异步逻辑的await,导致关闭流程没有正确执行。另外,stop_listening里直接调用wait_close()却没先关闭服务器,会导致它一直阻塞等待。

通用跨平台解决方案思路

因为Windows不支持Unix的信号处理机制,我们换个思路:通过捕获KeyboardInterrupt(Ctrl+C在所有平台都会触发这个异常)来触发优雅关闭,同时确保异步逻辑能正确执行。这里分两种场景给出方案:


方案1:单线程异步架构(推荐,更简洁)

如果你的服务不需要额外的同步线程任务,直接用asyncio.run托管整个异步生命周期是最省心的,跨平台兼容性拉满。

import asyncio
import functools
import websockets

PORT = 8765

class WebSocketPlayer:
    def __init__(self, id, websocket, server):
        self.id = id
        self.websocket = websocket
        self.server = server
    
    async def listen(self):
        # 保留你原来的消息处理逻辑
        try:
            while True:
                msg = await self.websocket.recv()
                print(f"Received from {self.id}: {msg}")
        except websockets.exceptions.ConnectionClosed:
            await self.server.remove_online_player(self.id)

class Server(object):
    def __init__(self):
        self.online_players = dict()
        self.online_players_lock = asyncio.Lock()
        self.websocket_server = None
    
    async def add_online_player(self, id, player):
        async with self.online_players_lock:
            self.online_players[id] = player
    
    async def remove_online_player(self, id):
        async with self.online_players_lock:
            self.online_players.pop(id, None)

async def on_connect(websocket, path, server):
    print("New user connected...")
    try:
        player_id = await websocket.recv()
        player = WebSocketPlayer(player_id, websocket, server)
        await server.add_online_player(player_id, player)
        await player.listen()
    except websockets.exceptions.ConnectionClosed:
        print("User disconnected unexpectedly")

async def main():
    server = Server()
    # 绑定处理函数并启动WebSocket服务器
    bound_handler = functools.partial(on_connect, server=server)
    server.websocket_server = await websockets.serve(
        bound_handler, "localhost", PORT, ping_timeout=None
    )
    print(f"Server running on ws://localhost:{PORT}")

    try:
        # 一直等待,直到捕获KeyboardInterrupt
        await asyncio.Event().wait()
    except KeyboardInterrupt:
        print("\nShutting down server gracefully...")

    # 优雅关闭流程
    # 1. 停止接受新连接
    server.websocket_server.close()
    # 2. 等待已有的连接全部关闭
    await server.websocket_server.wait_closed()
    # 3. 清理所有在线玩家的连接
    async with server.online_players_lock:
        for player in server.online_players.values():
            await player.websocket.close()
    print("Server shut down successfully")

if __name__ == "__main__":
    asyncio.run(main())

关键说明:

  • asyncio.run会自动管理事件循环的创建和销毁,不需要手动处理线程。
  • 用await asyncio.Event().wait()让服务一直运行,直到Ctrl+C触发KeyboardInterrupt。
  • 关闭流程严格按照"停新连接→等旧连接关闭→清理资源"的顺序,确保没有残留连接。

方案2:多线程架构(适配你原有的设计)

如果你必须用多线程分离asyncio loop和主线程,就要用asyncio.run_coroutine_threadsafe正确调度异步关闭任务,而非call_soon_threadsafe。

import asyncio
import functools
import threading
import time
import websockets

PORT = 8765

class WebSocketPlayer:
    def __init__(self, id, websocket, server):
        self.id = id
        self.websocket = websocket
        self.server = server
    
    async def listen(self):
        try:
            while True:
                msg = await self.websocket.recv()
                print(f"Received from {self.id}: {msg}")
        except websockets.exceptions.ConnectionClosed:
            await self.server.remove_online_player(self.id)

class Server(object):
    def __init__(self):
        self.online_players = dict()
        self.online_players_lock = asyncio.Lock()
        self.websocket_server = None
    
    async def add_online_player(self, id, player):
        async with self.online_players_lock:
            self.online_players[id] = player
    
    async def remove_online_player(self, id):
        async with self.online_players_lock:
            self.online_players.pop(id, None)
    
    def start(self):
        self.loop = asyncio.new_event_loop()
        # 设置守护线程,主线程退出时自动终止
        self.server_thread = threading.Thread(target=self._run_loop, args=(self.loop,))
        self.server_thread.daemon = True
        self.server_thread.start()

        try:
            # 主线程保持运行,等待Ctrl+C
            while True:
                time.sleep(1)
        except KeyboardInterrupt:
            print("\nShutting down server gracefully...")
            # 提交异步关闭任务到loop线程,并等待完成
            shutdown_future = asyncio.run_coroutine_threadsafe(
                self._stop_server(self.loop), self.loop
            )
            shutdown_future.result()  # 等待关闭流程完成
            self.server_thread.join()
    
    def _run_loop(self, loop):
        asyncio.set_event_loop(loop)
        bound_handler = functools.partial(on_connect, server=self)
        start_server_task = websockets.serve(
            bound_handler, "localhost", PORT, ping_timeout=None, loop=loop
        )
        self.websocket_server = loop.run_until_complete(start_server_task)
        print(f"Server running on ws://localhost:{PORT}")
        loop.run_forever()
    
    async def _stop_server(self, loop):
        # 1. 关闭WebSocket服务器
        self.websocket_server.close()
        await self.websocket_server.wait_closed()
        # 2. 清理所有玩家连接
        async with self.online_players_lock:
            for player in self.online_players.values():
                await player.websocket.close()
        # 3. 停止事件循环
        loop.stop()
        # 4. 关闭循环(可选,确保资源释放)
        loop.close()

async def on_connect(websocket, path, server):
    print("New user connected...")
    try:
        player_id = await websocket.recv()
        player = WebSocketPlayer(player_id, websocket, server)
        await server.add_online_player(player_id, player)
        await player.listen()
    except websockets.exceptions.ConnectionClosed:
        print("User disconnected unexpectedly")

if __name__ == "__main__":
    server = Server()
    server.start()

关键说明:

  • 用asyncio.run_coroutine_threadsafe提交异步关闭任务到loop线程,它返回concurrent.futures.Future,调用result()可等待异步任务完成。
  • 设置线程为守护线程,避免主线程退出后线程残留。
  • 关闭流程同样遵循"停新连接→等旧连接→清理资源→停loop"的顺序。

通用跨平台要点总结

  1. 捕获KeyboardInterrupt:这是跨平台处理Ctrl+C最可靠的方式,Windows和Unix系统都支持。
  2. 异步任务调度:如果用多线程,一定要用asyncio.run_coroutine_threadsafe处理异步函数,不能直接用call_soon_threadsafe。
  3. 优雅关闭顺序:先停止接受新连接,再等待已有连接关闭,最后清理资源,避免连接泄漏或资源未释放的问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 08:07:37