Pathway Web Dashboard:实时可视化监控流式数据管道的完整指南
本文讲解如何使用 Pathway Live Data Framework 内置的 Web Dashboard 对运行中的流式管道做实时监控:包括如何开启指标采集(pw.set_monitoring_config(detailed_metrics_dir=...))、如何启动 pathway web-dashboard 命令、各项参数与环境变量的含义,并结合开源仓库源码说明指标如何落盘为 SQLite 数据库、Dashboard 服务端如何读取这些数据库并向浏览器提供管道图(operator graph)、算子延迟与内存占用等视图。读完本文,你可以把该监控能力直接接入自己的管道部署,并从源码层面理解其工作机制。
概述与前提条件
内置 Web Dashboard 提供一套面向数据管道实时监控的图形化界面:交互式管道图绘制、指标图表与关键度量,可用于跟踪管道性能、检查数据流动状态并快速定位瓶颈。文档列出的核心能力包括:
- Advanced data flow inspection over the operator graph:在算子图(operator graph)上进行高级数据流检查,直观查看各算子之间数据的流动与状态;
- Latency tracking for each operator:对每个算子进行延迟跟踪,识别慢算子与堆积点;
- Memory usage tracking:内存占用跟踪,用于监控资源消耗。
许可证前提:Web Dashboard 需要 Pathway Live Data Framework 的 Scale 许可证(可免费申请)。集成测试 test_detailed_monitoring.py 验证了这一行为:在未设置许可证的情况下启用监控配置,run_all() 会抛出 api.EngineError,错误信息为 the feature(s) you used ["MONITORING"] require a Pathway license key, which is free.,说明监控功能(MONITORING 权益)由许可证校验强制把关。
快速上手
第一步:在管道中开启详细指标采集
在管道代码开头加入 pw.set_monitoring_config(detailed_metrics_dir="YOUR-DIR"),将许可证密钥与指标目录配置好后照常运行管道:
import pathway as pw
pw.set_license_key(key="YOUR-KEY")
pw.set_monitoring_config(detailed_metrics_dir=".")
# your pipeline here...
pw.run()
第二步:运行管道
按常规方式运行,例如 python your_pipeline.py,或使用 pathway spawn 启动。
提示:可通过环境变量
PATHWAY_METRICS_READER_INTERVAL_SECONDS配置指标读取频率,例如:PATHWAY_METRICS_READER_INTERVAL_SECONDS=1 python your_pipeline.py
第三步:打开 Dashboard
在另一个终端执行:
pathway web-dashboard
随后在浏览器中访问 Dashboard 即可看到管道图与指标图表。
源码级解析:配置是如何生效的
set_monitoring_config 的两个参数
该函数的完整签名定义在 config.py(set_monitoring_config,约 L161-L190):
def set_monitoring_config(
*,
server_endpoint: str | None = None,
detailed_metrics_dir: str | os.PathLike | None = None,
) -> None:
两个参数分别对应两类监控出口:
| 参数 | 作用 | 说明 |
|---|---|---|
server_endpoint |
监控服务端点 URL | 需要与 OTLP 兼容并支持 gRPC 协议,用于把指标/日志推送到 OpenTelemetry 兼容的接收端;传 None 会清除已有配置。 |
detailed_metrics_dir |
详细指标导出目录 | 将详细指标以 SQLite 数据库形式导出到该目录,Web Dashboard 正是读取这些数据库来渲染界面的;接受 str 或 os.PathLike,内部统一转为字符串。 |
函数把两个值写入全局的 PathwayConfig(一个 ContextVar 承载的 dataclass,定义于 config.py 的 PathwayConfig,L64-L114)。这些配置项同时可以通过环境变量提供,对应关系如下(均可从 PathwayConfig 的字段工厂函数中确认):
| 环境变量 | 对应配置 | 说明 |
|---|---|---|
PATHWAY_LICENSE_KEY |
license_key |
许可证密钥(Scale 许可证必填) |
PATHWAY_MONITORING_SERVER |
monitoring_server |
OTLP/gRPC 监控端点(空字符串时回退默认值) |
PATHWAY_DETAILED_METRICS_DIR |
detailed_metrics_dir |
详细指标 SQLite 导出目录 |
PATHWAY_METRICS_READER_INTERVAL_SECONDS |
metrics_reader_interval_secs |
指标读取间隔(秒),默认 None(使用内置默认值),传入时强制转换为 int |
也就是说,即使不写代码,也可以直接用环境变量把 PATHWAY_DETAILED_METRICS_DIR 指到某个目录来开启采集——集成测试 test_web_dashboard_spawn 正是用 env["PATHWAY_DETAILED_METRICS_DIR"] = str(metrics_dir) 这种方式驱动 Dashboard 的。
指标数据库的命名规则
测试 test_monitoring_license_detailed_metrics_created(test_detailed_monitoring.py L29-L44)验证了落盘行为:
- 测试先通过
monkeypatch.setenv("PATHWAY_RUN_ID", "test_run")设置运行 ID; - 配置
detailed_metrics_dir=metrics_path并执行run_all()后,断言目录存在,且其中生成了文件metrics_test_run.db。
即数据库文件命名为 metrics_<run_id>.db(run_id 由 PATHWAY_RUN_ID 环境变量决定),每次运行产出一个独立的 SQLite 文件,这也是 Dashboard 能按时间范围回放不同运行指标的基础。
pathway web-dashboard 命令的底层实现
CLI 命令实现位于 cli.py(L541-L565):
@cli.command()
@click.option(
"--detailed-metrics-dir",
type=str,
default=".",
help="directory in which metrics are stored",
)
@click.option(
"--port",
type=int,
default=8088,
)
def web_dashboard(detailed_metrics_dir, port):
...
env["PATHWAY_DETAILED_METRICS_DIR"] = detailed_metrics_dir
command = [
"uvicorn",
"pathway.web_dashboard.dashboard:app",
"--host", "0.0.0.0",
"--port", str(port),
]
subprocess.run(command, env=env)
从中可以确认几个关键事实:
- 参数:
--detailed-metrics-dir默认"."(当前目录),--port默认8088; - 启动方式:命令本质上是设置
PATHWAY_DETAILED_METRICS_DIR环境变量后,用 uvicorn 以0.0.0.0监听、启动pathway.web_dashboard.dashboard:app这个 ASGI 应用(FastAPI 风格的app对象),因此 Dashboard 是一个常驻的 Web 服务; - 启动时序:从 CHANGELOG.md 的记录看,
pathway web-dashboard会等待指标数据库创建完成而不是立即退出,即 Dashboard 可以先于/独立于管道启动,之后管道运行时写入的metrics_*.db会被读取。
测试 test_web_dashboard_spawn(L47-L88)进一步验证了服务端行为:先创建一个空的 metrics_empty.db 文件使 Dashboard 能启动,随后请求主页面确认返回 HTML,再请求 API 端点 /metrics/available_range,断言响应中包含 min 与 max 字段——即 Dashboard 通过该 API 向前端暴露当前指标库中可查询的时间范围,前端据此渲染图表。
典型目录与验证步骤
按上面的 Quickstart 走一遍后,可以按以下方式验证各部件:
- 检查指标目录:运行
pw.set_monitoring_config(detailed_metrics_dir=".")的管道后,当前目录应出现metrics_<run_id>.db的 SQLite 文件(可用任意 SQLite 客户端打开查看表结构); - 检查 Dashboard 服务:执行
pathway web-dashboard --detailed-metrics-dir . --port 8088,控制台输出Starting Pathway Live Data Framework Web Dashboard on port 8088...,浏览器访问http://<host>:8088/应加载出管道图界面; - 检查数据流/延迟/内存视图:在界面中查看算子图上的数据流标注、各算子延迟曲线与内存占用图表,与测试所覆盖的 API 能力(如
available_range时间范围查询)对应。
与 OTLP 监控通道的关系
set_monitoring_config 的两个参数对应两条互补的监控通道:
detailed_metrics_dir→ 本地 SQLite 落盘 → 内置 Web Dashboard 可视化(本文主题);server_endpoint→ 通过 OTLP/gRPC 推送到外部系统(如 OpenTelemetry Collector、Grafana 等)。
后者的完整配置流程(Collector 的 config.yaml、Grafana Loki 日志接入等)见同目录的姊妹文档 Monitoring a Pathway Live Data Framework Instance,其中结尾也明确指引读者到本文来了解 SQLite 导出 + 计算图可视化的做法。仓库中的监控示例项目 examples/projects/monitoring/monitoring_demo.py 与 examples/projects/monitoring/README.md 提供了可直接参考的完整监控配置演示。
小结
- Web Dashboard 是 Pathway 内置的实时管道监控界面,提供算子图数据流检查、逐算子延迟跟踪与内存占用监控三项核心能力,需要免费的 Scale 许可证;
- 开启方式是
pw.set_license_key(...)+pw.set_monitoring_config(detailed_metrics_dir=...),运行后指标以metrics_<run_id>.db命名的 SQLite 文件落盘,PATHWAY_METRICS_READER_INTERVAL_SECONDS控制读取频率; pathway web-dashboard(默认端口 8088)实际通过 uvicorn 启动pathway.web_dashboard.dashboard:app服务,读取上述 SQLite 数据库并经/metrics/available_range等 API 向前端提供可查询的时间范围;- 需要对接 OpenTelemetry/Grafana 等外部体系时,使用同一 API 的
server_endpoint参数,两条通道可并存使用。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0624
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00