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

如何在Flask同步方法中等待WebSocket异步响应?

在同步Flask接口中等待Asyncio/WebSocket异步响应的实现方案

完全可以实现类似Java CompletableFuture的功能,核心思路是用线程安全的未来对象在同步Flask线程和Asyncio事件循环之间传递响应结果,让同步接口阻塞等待异步回调返回的数据。

具体实现步骤及代码示例

1. 全局状态管理

首先维护一个线程安全的映射表,用来关联每个请求的唯一ID和对应的未来对象,同时保护映射表的线程安全:

from concurrent.futures import Future
from uuid import uuid4
from flask import Flask, request
import asyncio
import threading

app = Flask(__name__)

# 线程安全的请求ID与Future映射,存储待处理的请求
request_futures = {}
futures_lock = threading.Lock()

# 假设ws是已初始化完成的WebSocket客户端实例
ws = None

2. 同步Flask接口实现

在POST接口中生成唯一请求ID,发送WebSocket消息后阻塞等待Future的结果:

@app.route("/request", methods=["POST"])
def manage_request():
    data = request.get_json()
    # 生成唯一请求ID,用于关联请求和响应
    request_id = str(uuid4())
    data["request_id"] = request_id

    # 创建Future对象,用于接收异步响应
    future = Future()
    # 加锁保护映射表
    with futures_lock:
        request_futures[request_id] = future

    try:
        # 发送带请求ID的WebSocket消息
        ws.send(data)
        # 阻塞等待响应,可设置超时时间避免无限等待
        response = future.result(timeout=10)
        # 处理响应生成最终返回结果
        processed_response = make_something(response)
        return processed_response, 200
    except TimeoutError:
        return {"error": "请求超时,请重试"}, 408
    finally:
        # 清理映射表,避免内存泄漏
        with futures_lock:
            del request_futures[request_id]

def make_something(response):
    # 替换为你的实际业务处理逻辑
    return {"status": "success", "data": response}

3. 异步WebSocket消息回调修改

在on_message回调中解析请求ID,找到对应的Future并设置结果,唤醒等待的Flask线程:

async def on_message(message):
    # 解析消息数据(根据你的WebSocket消息格式调整)
    msg_data = message  # 假设消息已被解析为字典格式
    request_id = msg_data.get("request_id")

    # 处理无请求ID的消息(如广播、系统消息)
    if not request_id:
        await do_stuff(message)
        return

    # 加锁获取对应的Future对象
    with futures_lock:
        future = request_futures.get(request_id)
    
    # 如果Future存在且未完成,设置结果唤醒等待线程
    if future and not future.done():
        future.set_result(msg_data["response"])
    
    # 执行原有业务逻辑
    await do_stuff(message)

async def do_stuff(message):
    # 替换为你的原有异步业务逻辑
    pass

关键注意事项

  • 线程安全:必须用锁保护request_futures的读写操作,因为Flask是多线程处理请求,而Asyncio事件循环运行在单独线程,跨线程访问共享数据必须加锁。
  • 超时处理:future.result()一定要设置超时时间,避免接口因WebSocket无响应而一直阻塞。
  • 内存泄漏:在finally块中删除映射表中的Future对象,防止长期占用内存。
  • 请求ID唯一性:用UUID生成请求ID,确保每个请求的标识唯一,避免响应错乱。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 09:32:37