如何基于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
相关产品推荐
相关产品推荐

