首页
/ Pathway Live Data Framework 实例监控实战:OpenTelemetry Collector、Grafana Loki 接入与引擎遥测实现剖析

Pathway Live Data Framework 实例监控实战:OpenTelemetry Collector、Grafana Loki 接入与引擎遥测实现剖析

2026-09-04 17:17:36作者:申梦珏Efrain

本文基于 Pathway 仓库官方教程文档,讲解如何为一个 Pathway Live Data Framework(以下简称 Pathway LDF)流式管道搭建生产级监控:从用 Docker 启动 OpenTelemetry Collector(OTLP gRPC 接收端),到通过 pw.set_license_key / pw.set_monitoring_config 两行代码开启遥测导出,再到将日志推送至 Grafana Cloud Loki 并在 Explore 中检索验证。读完本篇,你将掌握完整的可复制的监控接入流程,并理解 Pathway 引擎内部究竟上报了哪些指标(进程 CPU/内存、输入/输出/算子延迟等)及其实现位置。

需要说明的前提:监控功能要求 Pathway Live Data Framework Scale 许可证(可免费获取)。未设置有效 license key 时,引擎会拒绝运行并抛出类似 the feature(s) you used ["MONITORING"] require a Pathway license key, which is free. 的错误——这一行为可以直接在 集成测试test_detailed_monitoring_insufficient_license 用例中得到印证。

一、监控架构总览:Pathway 通过 OTLP 协议外发遥测

对于业务关键的生产部署,必须保证当前的和历史上的瓶颈——延迟下降、吞吐波动、异常行为——都可被快速定位。Pathway LDF 采用的方案是 OpenTelemetry 协议(OTLP):引擎在运行时采集 traces、metrics、logs 三类遥测数据,通过 OTLP 的 gRPC 协议发送到外部接收端(如 OpenTelemetry Collector),再由 Collector 转发到任意可观测性后端(Grafana Loki、Prometheus、Tempo 等)。

从源码结构看,这一能力由 遥测模块 实现:它基于 opentelemetry / opentelemetry-otlp Rust 客户端,使用 PeriodicReader 周期性导出(源码中 PERIODIC_READER_INTERVAL 定义为 1 分钟,导出超时 3 秒),并通过 Resource 属性(如 service.namerun.idworker.id 等)标识每个运行实例。这也解释了后文在 Grafana 中用 service_name=~".*pathway" 过滤日志的原理。

二、配置并启动 OpenTelemetry Collector(Docker 方式)

Collector 有多种部署形态(Kubernetes Helm Chart、Grafana Alloy 等)。官方教程以最简单的 Docker + 官方 otel/opentelemetry-collector-contrib 镜像为例,前置要求是本地已安装 Docker。

2.1 编写 config.yaml:OTLP 接收器 + debug 导出器

先创建一个名为 config.yaml 的配置文件,启用 OTLP 接收器(gRPC 协议)与 debug 导出器:

receivers:
  otlp:
    protocols:
      grpc:

exporters:
  debug:
    verbosity: detailed

然后为 traces、metrics、logs 三类信号定义流水线,全部写入同一文件:

service:
  pipelines:
    traces:
      receivers: [otlp]
      exporters: [debug]
    metrics:
      receivers: [otlp]
      exporters: [debug]
    logs:
      receivers: [otlp]
      exporters: [debug]

此阶段的 debug 导出器仅用于验证数据通路——数据会在 Collector 的控制台日志中以 detailed 级别打印出来,无需任何后端账号即可确认 Pathway 已成功上报。

2.2 启动 Collector 并开放 gRPC 端口

docker run -v $(pwd)/config.yaml:/etc/otelcol-contrib/config.yaml -p <PORT>:4317 otel/opentelemetry-collector-contrib:latest

注意把 <PORT> 替换为你选择的宿主机端口。Collector 的 OTLP gRPC 接收端默认监听容器内 4317 端口,这就是后面 Pathway 端 server_endpoint 要指向的端口。

三、在 Pathway 管道中启用监控:两行 API 与底层配置机制

拿到 license key 后,只需把下面两行代码复制到管道脚本开头(pw.run() 之前)即可开启监控:

import pathway as pw

pw.set_license_key(key="YOUR-KEY")
pw.set_monitoring_config(server_endpoint="http://localhost:<PORT>")

# your pipeline here...

pw.run()

仓库中还提供了一个可直接运行的最小演示管道 monitoring_demo.py:它通过 pw.io.python.read 每秒产出一个整数并用 pw.reducers.sum 累加,适合用来冒烟验证监控数据流。运行方式:

python examples/projects/monitoring/monitoring_demo.py

3.1 源码层面的配置解析

两个 API 均定义在 config.py 中:

  • pw.set_license_key:把 key 写入当前 PathwayConfiglicense_key 字段,传 None 可清除;
  • pw.set_monitoring_config:签名关键字参数为 server_endpoint: str | Nonedetailed_metrics_dir: str | os.PathLike | None。文档字符串明确要求 endpoint 必须 OTLP 兼容并支持 gRPC 协议——这与第二节 Collector 只启用 grpc 接收端正好对应;传 None 可清除已有配置。

PathwayConfig 同时支持纯环境变量方式(对容器化部署更友好),从 dataclass 字段定义 可以读出完整的取值来源:

配置项 API 参数 环境变量 说明
许可证 set_license_key(key=...) PATHWAY_LICENSE_KEY 监控功能的前置条件
OTLP 端点 server_endpoint PATHWAY_MONITORING_SERVER OTLP gRPC 端点,如 http://localhost:4317
SQLite 明细指标目录 detailed_metrics_dir PATHWAY_DETAILED_METRICS_DIR 将指标另写本地 SQLite,供 Web Dashboard 可视化(见第六节)
指标读取间隔 PATHWAY_METRICS_READER_INTERVAL_SECONDS 控制指标读取节奏的可选整数

四、引擎实际上报了哪些数据

理解"Collector 里看到的东西"从何而来,有助于排查断连、延迟尖峰等问题。从 src/engine/telemetry.rs 中的常量定义,Pathway 至少导出以下指标:

指标名 含义
process.memory.usage 进程内存占用
process.cpu.utime / process.cpu.stime 进程用户态 / 内核态 CPU 时间
latency.input / latency.output 输入侧 / 输出侧延迟
operator.latency / operator.insertions / operator.deletions 各算子的延迟与增删行数

同时会附加 run.idworker.idroot.trace.idlicense.key 等资源属性。延迟的计算逻辑可在 dataflow/monitoring.rs 中找到:OperatorStats 持有每个算子的 time(时间戳)、lagdone 状态,latency(now) 以当前系统时间与事件时间戳的差值(毫秒)计算延迟,这正是 Grafana 里 latency.input 等曲线的来源。

启用 Collector 后运行管道,稍等片刻即可在 docker run 所在终端的 Collector 日志中看到 Pathway 上报的 OTLP 数据,说明链路打通。

五、接入 Grafana Cloud Loki:从调试模式到生产后端

OpenTelemetry Collector 的强项在于把同一份遥测扇出到不同后端。官方教程以 Grafana Loki 作为日志目的地(Grafana Cloud 免费套餐即可覆盖日志需求),并指出:完成日志接入后,把指标配置到 Grafana Prometheus、把 traces 配置到 Grafana Tempo 属于同构操作。

5.1 在 Collector 配置中增加 Loki 导出器

登录并创建 Grafana Cloud 账号后,在组织面板的 Loki 区块点击 "Send Logs" 按钮,页面会给出 Basic Auth 所需的用户名、Token 与推送 URL。请确保 Token 带有指标写入权限。将 教程文档 中的配置片段合并进 config.yaml

# 配置 basicauth 扩展(认证扩展)
extensions:
  basicauth/grafana_cloud_loki:
    client_auth:
      username:  <USER>
      password:  <TOKEN>

# 新增带认证的 loki 导出器
exporters:
  debug:
    verbosity: detailed
  loki/grafana_cloud_logs:
    endpoint: <URL>
    auth:
      authenticator: basicauth/grafana_cloud_loki

# 启用扩展,并把 loki 导出器加入 logs 流水线
service:
  extensions: [basicauth/grafana_cloud_loki]
  pipelines:
    logs:
      receivers: [otlp]
      exporters: [loki/grafana_cloud_logs, debug]
    # 其余 traces/metrics 流水线保持不变……

保存后重启 Collector 即可。仓库内的 examples/projects/monitoring/config.yaml 提供了"日志 + 指标 + traces 三路全通"的完整生产参考配置:为 Loki、Prometheus(prometheusremotewrite 导出器)、Tempo(otlp 导出器)分别配置了 basicauth 扩展与 batch 批处理处理器(send_batch_size: 10timeout: 30s),并借助 resource/loki 处理器插入 loki.format: raw 属性以控制日志渲染格式;配套的 docker-compose.yaml 通过 OTLP_GRPC_PORT(默认 4317)暴露端口,各凭据以环境变量注入。该目录下还附有可直接导入的 grafana-dashboard.json 看板,以及更完整的分步说明 README(要求 pip install -U pathway,版本 0.11.2 及以上)。

5.2 在 Grafana Explore 中验证日志

  1. 进入 Grafana 组织门户并打开 Grafana;
  2. 点击左侧 "Explore";
  3. 选择 grafanacloud-<your-org>-loki 数据源;
  4. 用 LogQL 过滤 service_name=~".*pathway"

此时应能看到 Pathway 引擎持续写入的运行日志,例如:

导入示例看板后,Grafana 中展示的资源使用曲线

六、不依赖外部系统的替代方案:本地 SQLite 明细指标

官方文档特别指出:除了把遥测导出到第三方系统,还可以把明细指标导出到本地 SQLite 数据库,并用 Web Dashboard 随时间可视化计算图(computation graph)。这正是 set_monitoring_config 的第二个参数 detailed_metrics_dir 的用途(对应环境变量 PATHWAY_DETAILED_METRICS_DIRCHANGELOG 中也记录了这一能力)。该模式的行为由 integration_tests/monitoring/test_detailed_monitoring.py 验证:设置 detailed_metrics_dir 后运行管道,目录下会生成形如 metrics_<run_id>.db 的 SQLite 文件;同一测试文件还验证了 pathway.web_dashboard.dashboard:app 可以以 uvicorn 启动并响应 /metrics/available_range 等 API。更多细节可参考仓库内的 Web Dashboard 教程

七、小结与延伸阅读

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
527
590
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
904
1.82 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
889
5.78 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.52 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.33 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
982
502
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384