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

如何基于GridDB时序数据库实现Streamlit仪表盘近实时更新?

GridDB + Streamlit 实时仪表盘高效集成方案

问题背景

我正在用Streamlit开发实时监控仪表盘,后端采用GridDB时序数据库,目标是实现每秒更新数据。当前通过griddb_python_client连接查询GridDB,依赖st.experimental_rerun()定期刷新仪表盘展示最新数据,但GridDB没有原生发布/订阅机制或事件触发器,只能靠轮询实现,效率较低。想寻求更智能、无需持续轮询的近实时/实时更新集成方式。

当前代码片段:

def get_latest_data():
    query = container.query("SELECT * WHERE timestamp > TIMESTAMPADD(SECOND, -5, NOW())")
    rs = query.fetch()
    return [rs.get_next() for _ in range(rs.size())]  

st.title("Real-time Dashboard")

data = get_latest_data()

st.write(data)

time.sleep(1)
st.experimental_rerun()

可行解决方案

方案1:优化轮询策略(轻量高效改进)

不用每秒全量拉取数据,通过以下方式减少无效查询和数据传输:

  • 记录上次查询的最大timestamp,每次仅查询该时间点之后的新数据
  • 动态调整轮询间隔:有新数据时缩短间隔(比如1秒),无数据时拉长间隔(比如3秒)
  • 复用GridDB连接,避免每次查询重新建立连接

示例代码:

import streamlit as st
import time
from griddb_python_client import GridStore

# 初始化会话状态,存储上次查询的最大时间戳
if "last_timestamp" not in st.session_state:
    st.session_state.last_timestamp = 0

# 复用GridDB连接(建议全局初始化,避免重复创建)
@st.cache_resource
def get_griddb_container():
    gridstore = GridStore.get_instance(your_connection_params)
    return gridstore.get_container("your_container_name")

container = get_griddb_container()

def get_latest_data(last_ts):
    if last_ts == 0:
        # 首次加载取最近5秒数据
        query = container.query("SELECT * WHERE timestamp > TIMESTAMPADD(SECOND, -5, NOW())")
    else:
        # 仅拉取上次之后的新数据
        query = container.query(f"SELECT * WHERE timestamp > {last_ts}")
    
    rs = query.fetch()
    data = [rs.get_next() for _ in range(rs.size())]
    if data:
        # 更新最新时间戳
        st.session_state.last_timestamp = max(item["timestamp"] for item in data)
    return data

st.title("Real-time Dashboard")

data = get_latest_data(st.session_state.last_timestamp)

if data:
    st.write(data)
else:
    st.write("暂无新数据")

# 动态调整刷新间隔
refresh_interval = 1 if data else 3
time.sleep(refresh_interval)
st.experimental_rerun()

方案2:引入中间消息队列(事件驱动)

借助消息队列实现“有新数据才查询”的逻辑,彻底避免无意义轮询:

  • 在数据写入GridDB的环节,同时向消息队列(比如Redis Pub/Sub、RabbitMQ)发送新数据通知
  • Streamlit端订阅消息队列,收到通知后再去GridDB拉取最新数据

数据写入端示例(伪代码):

def insert_data(data):
    # 写入GridDB
    container.put(data)
    # 发送更新通知到Redis
    import redis
    r = redis.Redis(host='localhost', port=6379, db=0)
    r.publish("griddb_data_updates", str(data["timestamp"]))

Streamlit端示例:

import streamlit as st
import redis
import time
from griddb_python_client import GridStore

# 初始化GridDB连接
@st.cache_resource
def get_griddb_container():
    gridstore = GridStore.get_instance(your_connection_params)
    return gridstore.get_container("your_container_name")

container = get_griddb_container()

# 初始化Redis订阅
r = redis.Redis(host='localhost', port=6379, db=0)
pubsub = r.pubsub()
pubsub.subscribe("griddb_data_updates")

# 会话状态初始化
if "last_timestamp" not in st.session_state:
    st.session_state.last_timestamp = 0
if "initial_data_loaded" not in st.session_state:
    # 首次加载历史数据
    query = container.query("SELECT * WHERE timestamp > TIMESTAMPADD(SECOND, -5, NOW())")
    rs = query.fetch()
    st.session_state.initial_data = [rs.get_next() for _ in range(rs.size())]
    st.write(st.session_state.initial_data)
    st.session_state.initial_data_loaded = True

st.title("Real-time Dashboard")

# 监听消息队列,有新通知时查询数据
while True:
    message = pubsub.get_message(timeout=1)
    if message and message["type"] == "message":
        # 拉取上次之后的新数据
        query = container.query(f"SELECT * WHERE timestamp > {st.session_state.last_timestamp}")
        rs = query.fetch()
        new_data = [rs.get_next() for _ in range(rs.size())]
        if new_data:
            st.write(new_data)
            st.session_state.last_timestamp = max(item["timestamp"] for item in new_data)
    time.sleep(0.1)

方案3:使用Streamlit异步组件(优化轮询阻塞)

如果GridDB Python客户端支持异步调用,结合Streamlit的异步能力,减少线程阻塞,提升轮询效率:

import streamlit as st
import asyncio
from griddb_python_client import GridStore

# 异步获取GridDB连接
@st.cache_resource
async def get_griddb_container_async():
    gridstore = await GridStore.get_instance_async(your_connection_params)
    return await gridstore.get_container_async("your_container_name")

async def get_latest_data_async(last_ts):
    container = await get_griddb_container_async()
    if last_ts == 0:
        query = await container.query_async("SELECT * WHERE timestamp > TIMESTAMPADD(SECOND, -5, NOW())")
    else:
        query = await container.query_async(f"SELECT * WHERE timestamp > {last_ts}")
    
    rs = await query.fetch_async()
    data = []
    while await rs.has_next():
        data.append(await rs.get_next_async())
    return data

st.title("Real-time Dashboard")

if "last_timestamp" not in st.session_state:
    st.session_state.last_timestamp = 0

async def main():
    data = await get_latest_data_async(st.session_state.last_timestamp)
    if data:
        st.write(data)
        st.session_state.last_timestamp = max(item["timestamp"] for item in data)
    await asyncio.sleep(1)
    st.experimental_rerun()

asyncio.run(main())

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 02:42:11