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

如何构建调用Python监控函数的Go API以实现WebSocket数据推送

可以实现,Python有对应机制模拟Go的goroutine+channel模式

完全可以在Python里实现类似Go的异步生产-消费模式,不用依赖gRPC,下面分两种适合高速度需求的场景说明:

一、异步协程方案(推荐,性能最优)

用Python的asyncio协程替代Go的goroutine,asyncio.Queue替代channel,搭配异步WebSocket库(比如websockets或aiohttp),这种方案开销极小,适合高并发和低延迟需求。

示例代码:

import asyncio
import websockets
from datetime import datetime

# 异步监控任务:生产数据并推入队列
async def monitor_task(data_queue):
    while True:
        # 替换为你的实际监控/格式化逻辑
        formatted_data = f"实时更新: {datetime.now().isoformat()}"
        await data_queue.put(formatted_data)
        await asyncio.sleep(0.5)  # 模拟数据获取间隔

# WebSocket处理协程:从队列取数据推送给客户端
async def websocket_consumer(websocket):
    data_queue = asyncio.Queue()
    # 启动监控协程(等价于Go的go关键字)
    asyncio.create_task(monitor_task(data_queue))
    
    while True:
        data = await data_queue.get()
        await websocket.send(data)
        data_queue.task_done()

# 启动WebSocket服务
async def main():
    async with websockets.serve(websocket_consumer, "0.0.0.0", 8765):
        await asyncio.Future()  # 保持服务运行

if __name__ == "__main__":
    asyncio.run(main())

关键说明:

  • asyncio.create_task用于启动后台协程,和goroutine一样是轻量级的,单进程可同时运行数千个协程
  • asyncio.Queue是异步安全的通道,自动处理协程间的同步,不会出现竞态问题
  • websockets库是纯异步实现,性能远高于同步WebSocket框架

二、同步线程方案(兼容阻塞式监控逻辑)

如果你的监控函数是阻塞式的(比如依赖不支持异步的第三方库),可以用threading.Thread启动监控线程,queue.Queue作为线程安全的通道,搭配支持线程的WebSocket框架(比如flask-socketio)。

示例代码:

import threading
import queue
import time
from datetime import datetime
from flask import Flask
from flask_socketio import SocketIO, emit

app = Flask(__name__)
# 开启线程模式,避免阻塞主线程
socketio = SocketIO(app, async_mode="threading")

data_queue = queue.Queue(maxsize=100)  # 限制队列长度,避免内存溢出

# 监控线程:生产数据
def monitor_thread():
    while True:
        formatted_data = f"实时更新: {datetime.now().isoformat()}"
        try:
            data_queue.put_nowait(formatted_data)
        except queue.Full:
            # 队列满时丢弃旧数据,可根据需求调整逻辑
            data_queue.get()
            data_queue.put(formatted_data)
        time.sleep(0.5)

# WebSocket连接处理:消费队列数据
@socketio.on("connect")
def handle_client_connect():
    while True:
        data = data_queue.get()
        emit("data_update", data)
        socketio.sleep(0)  # 让出CPU,避免阻塞其他连接

if __name__ == "__main__":
    # 启动后台监控线程
    threading.Thread(target=monitor_thread, daemon=True).start()
    socketio.run(app, host="0.0.0.0", port=8765)

关键说明:

  • queue.Queue是线程安全的,无需额外加锁就能在多线程间传递数据
  • 线程相比协程开销稍大,但对于一般的监控场景完全足够
  • 设置队列最大长度可以避免数据积压导致的内存问题

性能优化建议

  • 优先选择异步方案:协程的用户态调度没有线程切换的开销,在高并发场景下性能更接近Go的goroutine模式
  • 异步化IO操作:如果监控逻辑涉及文件读写、网络请求,一定要用异步IO库(比如aiofiles、aiohttp),避免阻塞事件循环
  • 避免不必要的序列化:如果数据是自定义格式,尽量减少JSON等序列化开销,直接传递二进制数据更高效

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 22:32:45