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

如何实现服务器Python守护线程向浏览器JS应用推送数据?

我太懂这种卡在最后一步的憋屈了——后端的队列消费、繁重计算都跑通了,结果就是没法把计算结果主动推给浏览器里的JS应用,尤其还是跨线程的场景,确实容易卡壳。你试过ZeroRPC的streams没成功,那咱们换几个更落地的方案试试,都是生产环境常用的:

方案1:用WebSocket实现双向实时推送(最通用)

WebSocket是浏览器原生支持的双向通信协议,完美适配“服务器主动推数据给客户端”的场景,Python端也有成熟的库支持。

Python端代码(结合守护线程+队列)

先安装依赖:pip install websockets

import asyncio
import websockets
from queue import Queue
import threading
import time

# 全局任务队列,守护线程从这里取数据计算
task_queue = Queue()
# 保存所有活跃的WebSocket客户端连接
active_connections = set()

# 处理新的WebSocket连接
async def register_client(websocket):
    active_connections.add(websocket)
    try:
        await websocket.wait_closed()  # 保持连接直到客户端断开
    finally:
        active_connections.remove(websocket)

# 给所有客户端推送计算结果
async def push_result_to_clients(result):
    if active_connections:
        # 可以把结果转成JSON格式,方便前端解析
        message = f"{{'type': 'compute_result', 'data': {result}}}"
        # 批量推送给所有在线客户端
        await asyncio.gather(*[ws.send(message) for ws in active_connections])

# 守护线程:消费队列+执行繁重计算
def compute_worker_thread():
    while True:
        task_data = task_queue.get()
        if task_data is None:  # 用None作为退出信号
            break
        
        # 模拟繁重计算(替换成你的实际业务逻辑)
        time.sleep(3)
        compute_result = task_data * 2
        
        # 关键:在非异步线程中调用异步推送函数
        asyncio.run_coroutine_threadsafe(
            push_result_to_clients(compute_result), 
            asyncio.get_event_loop()
        )
        task_queue.task_done()

if __name__ == "__main__":
    # 启动WebSocket服务器
    start_server = websockets.serve(register_client, "localhost", 8765)
    loop = asyncio.get_event_loop()
    loop.run_until_complete(start_server)

    # 启动计算守护线程
    compute_thread = threading.Thread(target=compute_worker_thread, daemon=True)
    compute_thread.start()

    # 模拟往队列中添加任务(实际场景中可以从其他线程/进程喂数据)
    def feed_tasks():
        for i in range(5):
            task_queue.put(i)
            time.sleep(1)
        task_queue.put(None)  # 发送退出信号
    
    feed_thread = threading.Thread(target=feed_tasks)
    feed_thread.start()

    # 启动事件循环
    try:
        loop.run_forever()
    except KeyboardInterrupt:
        pass
    finally:
        loop.close()

前端JS代码

// 建立WebSocket连接
const ws = new WebSocket('ws://localhost:8765');

// 接收后端推送的结果
ws.onmessage = function(event) {
    const result = JSON.parse(event.data);
    console.log('收到计算结果:', result.data);
    // 这里可以把结果渲染到页面DOM中
    document.getElementById('result-container').innerText += result.data + '\n';
};

// 处理连接断开,自动重连
ws.onclose = function() {
    console.log('WebSocket连接断开,3秒后尝试重连...');
    setTimeout(() => window.location.reload(), 3000);
};
方案2:修复ZeroRPC Streams的用法

如果你坚持用ZeroRPC,可能之前的stream实现没踩对生成器的点——ZeroRPC的stream是通过Python生成器持续返回结果,JS端需要通过回调的more参数判断是否还有后续数据。

Python端代码

import zerorpc
from queue import Queue
import threading
import time

task_queue = Queue()

class ComputeService:
    def stream_compute_results(self):
        """用生成器返回持续的计算结果"""
        while True:
            task_data = task_queue.get()
            if task_data is None:
                break
            # 模拟繁重计算
            time.sleep(2)
            yield task_data * 2
            task_queue.task_done()

def task_feeder():
    # 模拟喂任务
    for i in range(5):
        task_queue.put(i)
        time.sleep(1)
    task_queue.put(None)

if __name__ == "__main__":
    server = zerorpc.Server(ComputeService())
    server.bind("tcp://0.0.0.0:4242")
    
    # 启动喂任务线程
    feeder_thread = threading.Thread(target=task_feeder)
    feeder_thread.start()
    
    server.run()

前端JS代码(依赖zerorpc-js)

先安装:npm install zerorpc

const zerorpc = require("zerorpc");
const client = new zerorpc.Client();
client.connect("tcp://localhost:4242");

// 调用stream方法,通过回调接收持续结果
client.invoke("stream_compute_results", (error, res, more) => {
    if (error) {
        console.error('ZeroRPC出错:', error);
        return;
    }
    console.log('收到计算结果:', res);
    // more为false时表示流结束,可以关闭连接
    if (!more) {
        client.close();
    }
});
方案3:Server-Sent Events (SSE) 单向推送(极简方案)

如果你的场景只需要服务器推数据给前端,不需要前端发消息给后端,SSE是更轻量的选择——浏览器原生支持,不需要额外JS库,Python端用Flask就能快速实现。

Python端代码(Flask)

from flask import Flask, Response
from queue import Queue
import threading
import time

app = Flask(__name__)
task_queue = Queue()
result_queue = Queue()

def compute_worker():
    while True:
        task_data = task_queue.get()
        if task_data is None:
            break
        time.sleep(2)
        result_queue.put(task_data * 2)
        task_queue.task_done()

# 生成SSE事件流
def generate_sse_events():
    while True:
        result = result_queue.get()
        # SSE的格式要求:data: 内容\n\n
        yield f"data: {result}\n\n"
        result_queue.task_done()

@app.route('/stream-results')
def stream_results():
    return Response(generate_sse_events(), mimetype='text/event-stream')

def feed_tasks():
    for i in range(5):
        task_queue.put(i)
        time.sleep(1)
    task_queue.put(None)

if __name__ == '__main__':
    compute_thread = threading.Thread(target=compute_worker, daemon=True)
    compute_thread.start()
    feed_thread = threading.Thread(target=feed_tasks)
    feed_thread.start()
    app.run(debug=True, threaded=True)

前端JS代码

// 建立SSE连接
const evtSource = new EventSource('/stream-results');

// 接收推送结果
evtSource.onmessage = function(event) {
    console.log('收到SSE推送:', event.data);
    document.getElementById('result-box').innerText += event.data + '\n';
};

// 处理连接错误
evtSource.onerror = function() {
    console.log('SSE连接出错,关闭连接');
    evtSource.close();
};

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:36:20