首页
/ Pathway Web Dashboard:实时可视化监控流式数据管道的完整指南

Pathway Web Dashboard:实时可视化监控流式数据管道的完整指南

2026-09-04 21:05:48作者:裘晴惠Vivianne

本文讲解如何使用 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.pyset_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 正是读取这些数据库来渲染界面的;接受 stros.PathLike,内部统一转为字符串。

函数把两个值写入全局的 PathwayConfig(一个 ContextVar 承载的 dataclass,定义于 config.pyPathwayConfig,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_createdtest_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>.dbrun_idPATHWAY_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 环境变量后,用 uvicorn0.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,断言响应中包含 minmax 字段——即 Dashboard 通过该 API 向前端暴露当前指标库中可查询的时间范围,前端据此渲染图表。

典型目录与验证步骤

按上面的 Quickstart 走一遍后,可以按以下方式验证各部件:

  1. 检查指标目录:运行 pw.set_monitoring_config(detailed_metrics_dir=".") 的管道后,当前目录应出现 metrics_<run_id>.db 的 SQLite 文件(可用任意 SQLite 客户端打开查看表结构);
  2. 检查 Dashboard 服务:执行 pathway web-dashboard --detailed-metrics-dir . --port 8088,控制台输出 Starting Pathway Live Data Framework Web Dashboard on port 8088...,浏览器访问 http://<host>:8088/ 应加载出管道图界面;
  3. 检查数据流/延迟/内存视图:在界面中查看算子图上的数据流标注、各算子延迟曲线与内存占用图表,与测试所覆盖的 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.pyexamples/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 参数,两条通道可并存使用。
登录后查看全文
热门项目推荐
相关项目推荐