Bokeh Server多会话共享实时数据的实现方案咨询
我刚接触Bokeh,但Python用得很熟。我想做一个基于Bokeh Server的实时数据可视化 viewer,需求是用单独线程获取数据,然后把数据同步给所有打开的文档/会话。有没有靠谱的实现方案?
我目前的尝试是:服务器启动时就启动数据线程,新数据一来就遍历所有会话更新图表。单会话时没问题,但打开第二个会话就报错了(错误栈在下面),而且页面加载后只有初始数据,要等新数据到才会更新。
我的尝试代码
app_hooks.py
import time import threading import bokeh.server.contexts from bokeh.plotting import Document from bokeh.models import ColumnDataSource keep_running = True data = { # some initial data 'x': [1, 2, 3, 4, 5], 'y': [6, 7, 2, 4, 7] } def callback(document: Document, new_data: dict): model = document.get_model_by_name("line") if model is None: print("Model not found in document") return assert isinstance(model, bokeh.models.renderers.glyph_renderer.GlyphRenderer) source: ColumnDataSource = model.data_source # source.stream(new_data, rollover=100) source.data = new_data def retrieve_data(i: int) -> dict[str, list[float]]: # idea: use a PUB/SUB pattern to retrieve data global data time.sleep(1) # it might take some time until new data is available. Can be seconds to hours. data['x'].append(i) data['y'].append(i * 0.5 % 5) return data def on_server_loaded(server_context: bokeh.server.contexts.BokehServerContext): # If present, this function executes when the server starts. def add_data_worker(): global keep_running i=0 while keep_running: print(f"Iteration {i + 1}") new_data = retrieve_data(i) for session in server_context.sessions: print(session.destroyed, session.expiration_requested) session.document.add_next_tick_callback(lambda: callback(session.document, new_data)) i += 1 thread = threading.Thread(target=add_data_worker, daemon=True) thread.start() def on_server_unloaded(server_context: bokeh.server.contexts.BokehServerContext): # If present, this function executes when the server shuts down. global keep_running keep_running = False
main.py
from bokeh.plotting import figure, curdoc from bokeh.models import ColumnDataSource from .app_hooks import data source = ColumnDataSource(data=data) p = figure(title="Simple Line Example", x_axis_label='x', y_axis_label='y') p.line(x="x", y="y", legend_label="My Value", line_width=2, source=source, name="line") curdoc().add_root(p)
报错信息
2025-06-23 08:43:14,863 Exception in callback functools.partial(<bound method IOLoop._discard_future_result of <tornado.platform.asyncio.AsyncIOMainLoop object at 0x0000029D6617B1F0>>, <Task finished name='Task-204' coro=<Server Session.with_document_locked() done, defined at C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\bokeh\server\session.py:77> exception=RuntimeError('_pending_writes should be non-None when we have a document lock, and we should have the lock when the document changes')>) Traceback (most recent call last): File "C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\tornado\ioloop.py", line 740, in _run_callback ret = callback() File "C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\tornado\ioloop.py", line 764, in _discard_future_result future.result() File "C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\bokeh\server\session.py", line 94, in _needs_document_lock_wrapper result = func(self, *args, **kwargs) File "C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\bokeh\server\session.py", line 226, in with_document_locked return func(*args, **kwargs) File "C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\bokeh\document\callbacks.py", line 495, in wrapper return invoke_with_curdoc(doc, invoke) File "C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\bokeh\document\callbacks.py", line 453, in invoke_with_curdoc return f() File "C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\bokeh\document\callbacks.py", line 494, in invoke return f(*args, **kwargs) File "C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\bokeh\document\callbacks.py", line 184, in remove_then_invoke return callback() File "C:\XtProgramFiles\python_libs_v3\dev\bokeh_tests\app_hooks.py", line 42, in <lambda> session.document.add_next_tick_callback(lambda: callback(session.document, new_data)) File "C:\XtProgramFiles\python_libs_v3\dev\bokeh_tests\app_hooks.py", line 22, in callback source.data = new_data File "C:\XtProgramFiles\Miniconda3\envs\dev3\lib\site-packages\bokeh\core\has_props.py", line 336, in __setattr__ return super().__setattr__(name,
解决方案分析与优化
你的思路方向完全没问题,问题主要出在闭包变量绑定和会话有效性判断上,咱们一步步来修正:
1. 解决闭包延迟绑定的坑
你用lambda: callback(session.document, new_data)时,lambda不会立刻捕获当前的session和new_data,而是在执行时才去取变量的最新值——这就导致所有回调最终都用循环最后一次的会话和数据,直接引发多会话冲突。
修改成用默认参数提前绑定当前值:
session.document.add_next_tick_callback(lambda doc=session.document, data=new_data: callback(doc, data))
2. 过滤无效会话
遍历server_context.sessions时,可能会包含已经销毁或标记为过期的会话,操作这些会话的文档肯定会报错。加上过滤条件:
for session in server_context.sessions: if not session.destroyed and not session.expiration_requested: session.document.add_next_tick_callback(lambda doc=session.document, data=new_data: callback(doc, data))
3. 解决初始数据同步问题
新会话加载后只有初始数据,是因为没有主动把当前最新数据同步给它。咱们加一个on_session_created钩子,在新会话建立时直接同步最新数据:
在app_hooks.py里添加:
def on_session_created(session_context: bokeh.server.contexts.BokehServerSessionContext): # 新会话创建时立刻同步当前最新数据 doc = session_context.session.document # 复制一份数据,避免后续线程修改影响当前回调 latest_data = {k: v.copy() for k, v in data.items()} doc.add_next_tick_callback(lambda d=doc, nd=latest_data: callback(d, nd))
4. 更高效的数据源更新方式
你注释掉的source.stream()其实是Bokeh推荐的实时数据更新方式,它只会追加新数据并自动处理滚动(通过rollover控制保留的最大数据量),比直接替换source.data更高效,还能减少前端重绘的开销。把回调里的代码改回:
source.stream(new_data, rollover=100)
优化后的核心逻辑
现在整个流程就通顺了:
- 服务器启动后,后台线程持续获取数据并更新全局
data字典 - 新数据到来时,遍历所有有效会话,通过
add_next_tick_callback在Tornado IO线程安全更新每个会话的图表 - 新用户打开页面时,
on_session_created会立刻把最新数据同步过去,不用等下一次数据更新
这样就能完美支持多会话的实时数据共享,而且不会再出现线程安全相关的错误。
内容来源于stack exchange

