如何从异步函数调用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
相关产品推荐
相关产品推荐

