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

如何从异步函数调用websocket.send()并实现跨异步函数数据传递

解决POST数据转发至WebSocket客户端的问题

我明白你现在的困境——要把第三方POST过来的数据实时推给WebSocket客户端,用Sanic和Rx探索的时候卡壳了对吧?核心问题其实是要建立一个消息广播通道,把HTTP接口收到的数据转发给所有在线的WebSocket连接。我来给你梳理两种可行的方案,先从最直观的Sanic原生实现入手,再整合Rx的思路解决你遇到的参数问题。

方案一:Sanic原生实现(无需Rx)

这个思路很直接:维护一个全局的WebSocket连接集合,当POST接口收到数据时,遍历集合给每个客户端发消息。

from sanic import Sanic
from sanic.response import json
from sanic.websocket import WebSocketProtocol

app = Sanic("post_to_websocket")

# 保存所有活跃的WebSocket连接
active_ws_connections = set()

@app.websocket('/ws')
async def websocket_handler(request, ws):
    # 连接建立时,把当前ws加入集合
    active_ws_connections.add(ws)
    try:
        # 保持连接存活(不需要处理客户端消息的话,这里可以空循环)
        while True:
            await ws.recv()  # 或者用await asyncio.sleep(3600)
    finally:
        # 连接断开时,从集合移除
        active_ws_connections.remove(ws)

@app.post('/webhook')
async def handle_post(request):
    # 解析POST过来的数据
    post_data = request.json or request.form
    if not post_data:
        return json({"error": "No data received"}, status=400)
    
    # 把数据转发给所有WebSocket客户端
    for ws in active_ws_connections.copy():  # 用copy避免遍历中集合变化报错
        try:
            await ws.send(str(post_data))  # 根据需要序列化数据,比如json.dumps
        except Exception as e:
            print(f"Failed to send to client: {e}")
            # 移除失效的连接
            active_ws_connections.discard(ws)
    
    return json({"status": "success", "sent_to": len(active_ws_connections)})

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=8000, protocol=WebSocketProtocol)

关键点说明

  • active_ws_connections用set存储,方便快速添加/移除连接
  • WebSocket handler里用try...finally确保连接断开时能从集合中移除,避免内存泄漏
  • POST接口里用copy()遍历集合,防止遍历过程中集合被修改(比如客户端突然断开)导致报错
  • 发送消息时捕获异常,及时处理已经失效的连接

方案二:整合RxPy实现响应式广播

你想用Rx的话,核心是用Subject——它既是Observer(可以接收事件)又是Observable(可以被订阅),刚好适合做消息广播,能解决你之前遇到的observable_message()需要参数的问题。

from sanic import Sanic
from sanic.response import json
from sanic.websocket import WebSocketProtocol
from rx import Subject
import asyncio

app = Sanic("rx_post_to_websocket")

# 创建一个Subject作为消息广播中心
message_bus = Subject()

@app.websocket('/ws')
async def websocket_handler(request, ws):
    # 订阅message_bus,收到消息就异步发给客户端
    subscription = message_bus.subscribe(
        on_next=lambda data: asyncio.create_task(ws.send(str(data))),
        on_error=lambda e: print(f"Error in subscription: {e}"),
        on_completed=lambda: print("Subscription completed")
    )
    try:
        # 保持连接存活
        while True:
            await ws.recv()
    finally:
        # 连接断开时取消订阅,释放资源
        subscription.dispose()

@app.post('/webhook')
async def handle_post(request):
    post_data = request.json or request.form
    if not post_data:
        return json({"error": "No data received"}, status=400)
    
    # 把POST数据推送至消息总线,所有订阅的WebSocket客户端都会收到
    message_bus.on_next(post_data)
    
    return json({"status": "success"})

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=8000, protocol=WebSocketProtocol)

关键点说明

  • message_bus = Subject()是核心枢纽,所有POST数据通过message_bus.on_next(post_data)推送,无需手动给Observable传参数
  • 每个WebSocket连接建立时自动订阅消息总线,用asyncio.create_task把同步的Rx回调包装成异步任务,适配Sanic的异步环境
  • 连接断开时调用subscription.dispose()取消订阅,避免无效订阅占用资源

额外注意事项

  • 数据序列化:如果POST数据是复杂JSON结构,建议用json.dumps(post_data)转换成字符串再发送给WebSocket客户端,避免格式错误
  • 多进程适配:如果你的Sanic启用多进程模式,原生方案的连接集合或Rx的Subject需要用Redis等共享存储实现跨进程同步,单进程模式下无需额外处理
  • 异常兜底:无论哪种方案,都要做好发送失败的异常处理,避免单个客户端的连接问题影响整个广播流程

内容的提问来源于stack exchange,提问作者C. Myles

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 03:34:41