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

资源预订系统中,延迟通知等待客户端的最优方案咨询

资源预订系统的实现方案与异步编程优化建议

需求背景

我正在开发一套包含服务端与客户端的资源预订系统,客户端为简易命令行工具,可完成资源的增删及预订操作。核心需求包括:

  • 客户端使用完资源后需告知服务端,资源将重新变为可用状态;
  • 若预订时资源不可用,客户端会保持等待,服务端需在资源可用时通知等待的客户端。

当前实现思路

目前的方案是:客户端发起预订请求时启动一个小型服务,并将自身地址与端口包含在请求中,服务端据此向客户端发送通知。为简化演示,我编写了一个小示例,客户端将自身地址发送给服务端,服务端定期向客户端发送消息,直到客户端无响应为止。

示例代码

server.py

import asyncio
from typing import List

from aiohttp import ClientConnectorError, ClientSession, web
from aiohttp.web_request import Request


async def handle(request: Request):
    subscriber = await request.text()

    print(f"New subscriber {subscriber}")

    request.app["subscribers"].append(subscriber)

    return web.Response()


app = web.Application()
app.add_routes([web.post("/", handle)])


async def periodic_publish(subscribers: List[str]):
    while True:
        await asyncio.sleep(2)

        print(f"{subscribers=}")

        async with ClientSession() as session:
            to_remove = []
            for i, subscriber in enumerate(subscribers):
                try:
                    async with session.post(
                        subscriber, data="Here is your subscription"
                    ) as response:
                        body = await response.text()
                        print(f"Response from subscriber: {body}")
                except ClientConnectorError:
                    print(f"Could not connect to {subscriber}")
                    to_remove.append(i)

        for subscriber_i in reversed(to_remove):
            subscribers.pop(subscriber_i)


if __name__ == "__main__":
    subscribers = []
    loop = asyncio.get_event_loop()
    loop.create_task(periodic_publish(subscribers))
    app["subscribers"] = subscribers
    web.run_app(app, port=8080, loop=loop)

client.py

import asyncio
import sys

import aioconsole
import aiohttp
from aiohttp import web
from aiohttp.web_request import Request

port = int(sys.argv[1])


async def handle(request: Request):
    message = await request.text()

    print(f'New message "{message}"')

    return web.Response(text="Thank you!")


app = web.Application()
app.add_routes([web.post("/", handle)])


async def cli():
    async with aiohttp.ClientSession() as session:
        async with session.post(
            "http://localhost:8080/", data=f"http://localhost:{port}/"
        ) as response:
            print("Status:", response.status)
            print("Content-type:", response.headers["content-type"])

            body = await response.text()
            print("Body:", body)

    await aioconsole.aprint("Hit enter to exit.")
    await aioconsole.ainput()
    await aioconsole.aprint("Bye!")


loop = asyncio.get_event_loop()
runner = web.AppRunner(app)
loop.run_until_complete(runner.setup())
site = web.TCPSite(runner, port=port)
loop.run_until_complete(site.start())
loop.run_until_complete(cli())

替代实现方式

考虑到客户端是可被用户随时终止的命令行工具,除了当前的客户端启动服务端回调的方式,还有以下几种方案:

  • 长轮询(Long Polling):客户端发起预订请求后,服务端保持连接打开,直到资源可用或超时。若超时,客户端重新发起请求。无需客户端启动额外服务,实现简单,适合命令行场景:

    1. 客户端发起POST /book-resource请求,携带资源ID;
    2. 服务端检查资源状态,可用则直接返回成功;不可用则挂起请求,直到资源被释放;
    3. 客户端收到响应后使用资源,完毕后调用POST /release-resource通知服务端。
  • WebSocket 连接:客户端与服务端建立长连接,所有操作通过该连接完成,实时性高且无需客户端暴露端口:

    1. 客户端启动后建立WebSocket连接;
    2. 发送预订请求,服务端记录该客户端连接;
    3. 资源可用时,服务端通过WebSocket发送通知;
    4. 客户端使用完资源后,通过WebSocket发送释放消息;
    5. 客户端退出时主动关闭连接,服务端清理等待记录。
  • 基于消息队列的异步通知:引入轻量级消息队列(如Redis Pub/Sub),客户端订阅资源通知频道:

    1. 客户端发起预订请求时,订阅对应资源的可用通知频道;
    2. 服务端释放资源时,向对应频道发送消息;
    3. 客户端收到消息后获取资源,使用完毕通知服务端;
    4. 客户端退出时自动取消订阅,消息队列处理未消费消息。

asyncio 和 aiohttp 用法建议

针对首次使用异步编程的代码,给出以下优化点:

  • 避免手动操作事件循环:Python 3.7+推荐使用asyncio.run()替代手动管理循环,代码更简洁规范。例如server.py的主函数可重构为:

    if __name__ == "__main__":
        subscribers = []
        app["subscribers"] = subscribers
        async def main():
            asyncio.create_task(periodic_publish(subscribers))
            await web.run_app(app, port=8080)
        asyncio.run(main())
    
  • 线程安全的订阅者管理:当前subscribers列表在多任务中读写存在并发风险,建议用asyncio.Lock保护:

    from asyncio import Lock
    
    subscribers_lock = Lock()
    
    async def handle(request: Request):
        subscriber = await request.text()
        print(f"New subscriber {subscriber}")
        async with subscribers_lock:
            request.app["subscribers"].append(subscriber)
        return web.Response()
    
  • 客户端服务优雅关闭:当前客户端未关闭启动的web服务,可在cli()函数末尾添加关闭逻辑:

    async def cli():
        # 原请求代码...
        await aioconsole.aprint("Bye!")
        await runner.cleanup()
    
  • 增强错误处理:除ClientConnectorError外,可捕获更多网络异常(如超时、连接重置),同时给客户端返回明确的错误状态码,便于调试。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 02:32:51