ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Pathway 实时数据处理监控实战:使用 OpenTelemetry Collector 与 Grafana Cloud 构建可观测性

2026/9/8 23:29:29 拓冰建站 浏览量
Pathway 实时数据处理监控实战:使用 OpenTelemetry Collector 与 Grafana Cloud 构建可观测性 Pathway 实时数据处理监控实战使用 OpenTelemetry Collector 与 Grafana Cloud 构建可观测性【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway本指南以 Pathway 仓库中的监控示例为主体完整讲解如何用 Docker Compose 部署一个 OpenTelemetry Collector把 Pathway 管线产生的指标Metrics、日志Logs与链路追踪Traces推送到 Grafana Cloud并通过官方 Grafana Dashboard 可视化资源占用与端到端延迟。读完本文你将掌握从 Pathway 侧开启监控上报、到 Collector 侧完成 OTLP 接收与多后端转发、再到 Grafana 侧导入仪表盘的完整落地路径。该主题在仓库中有两份互为镜像的说明docs/2.developers/7.templates/ETL/_readmes/monitoring.md 与其对应模板目录下的 examples/projects/monitoring/README.md所有配套文件均位于 examples/projects/monitoring/。一、方案组成与整体链路整套监控方案围绕三类开源标准与技术栈展开仓库中的文件与之一一对应角色技术仓库文件Pathway 管线OTLP gRPC 上报端monitoring_demo.py遥测收集器OpenTelemetry Collectorcontrib 镜像config.yaml、docker-compose.yaml指标后端Grafana Cloud Prometheusremote writeconfig.yaml 中的 prometheusremotewrite exporter日志后端Grafana Cloud Loki同上loki exporter链路后端Grafana Cloud Tempo同上otlp exporter可视化Grafana Cloud Dashboardgrafana-dashboard.json数据流向可以概括为一条三段式链路Pathway 侧调用pw.set_monitoring_config(server_endpoint...)后管线在运行时通过 OTLP/gRPC 协议把日志与遥测数据发往 Collector 暴露的 4317 端口Collector 侧otel/opentelemetry-collector-contrib作为统一入口接收三类遥测数据经 batch 等处理器聚合后按管线分别转发到 Grafana Cloud 的 Prometheus、Loki 与 TempoGrafana Cloud 侧导入仓库附带的 dashboard JSON即可在面板上查看输入/输出延迟、进程 CPU 与内存占用等指标并联动查看日志与 trace。从源码层面可以印证这一设计。Pathway 的遥测模块位于 python/pathway/internals/graph_runner/telemetry.py其模块注释明确说明启用后遥测数据管线元数据与进程指标会通过 OpenTelemetry 协议发送到指定端点转发到监控服务器的数据同时包含日志与遥测数据而发送到 Pathway 官方遥测服务器的数据不包含日志。同一文件还导入了OTLPLogExporter与OTLPSpanExporter分别负责日志与 trace 的 gRPC 导出说明 Pathwalker 的监控端点要求是 OTLP 兼容且支持 gRPC 协议的服务默认 telemetry 与 monitoring 均处于关闭状态。二、环境准备Prerequisites在动手前需要准备以下几项Docker 与 Docker Compose用于以容器方式拉起 OpenTelemetry Collector安装方式参考 Docker 官方文档Grafana Cloud 免费账号用于申请指标、日志与链路三个服务的写入凭据URL、用户名、密码/令牌Pathway 0.11.2 及以上版本监控示例基于该版本及以上的 API 编写可通过 pip 安装或升级pip install -U pathwayPathway 的监控能力需要有效的 License Key。在 python/pathway/internals/config.py 中set_monitoring_config的 docstring 明确指出它 Requires a valid Pathway Live Data Framework Scale license key因此必须在管线代码中先调用pw.set_license_key(key...)该函数定义于 config.py可前往 Pathway 官网免费申请。三、OpenTelemetry Collector 配置详解Collector 的行为由 config.yaml 全权控制官方称之为 pipeline 配置由 receivers、processors、exporters、extensions 与 service 五部分构成。3.1 receiversOTLP gRPC 监听入口receivers: otlp: protocols: grpc: endpoint: 0.0.0.0:4317Collector 在宿主机所有网卡的 4317 端口上监听 OTLP/gRPC 请求这正是 Pathway 侧server_endpointhttp://localhost:port需要指向的地址与端口4317 是 OTLP gRPC 的官方默认端口与下文的OTLP_GRPC_PORT默认值保持一致。3.2 extensions面向 Grafana Cloud 三后端的 BasicAuth 鉴权extensions: basicauth/grafana_cloud_tempo: client_auth: username: ${env:TEMPO_USERNAME} password: ${env:TEMPO_PASSWORD} basicauth/grafana_cloud_prometheus: client_auth: username: ${env:PROMETHEUS_USERNAME} password: ${env:PROMETHEUS_PASSWORD} basicauth/grafana_cloud_loki: client_auth: username: ${env:LOKI_USERNAME} password: ${env:LOKI_PASSWORD}Tempo、Prometheus、Loki 三套后端各配置一个独立的 basicauth extension用户名与密码通过环境变量注入避免把凭据硬编码进配置文件。这三个 extension 随后会在 service.extensions 中被启用并被各 exporter 通过auth.authenticator字段引用。3.3 processors批量发送与日志格式processors: batch: send_batch_size: 10 timeout: 30s resource/loki: attributes: - action: insert key: loki.format value: rawbatch对遥测数据做攒批攒满 10 条或等待 30 秒后一次性下发减少与远端后端之间的网络往返resource/loki为发往 Loki 的日志资源插入loki.format: raw属性告知 Loki 以 raw 文本方式解析日志行。3.4 exporters调试与三后端转发exporters: debug: verbosity: detailed otlp/grafana_cloud_traces: endpoint: ${env:TEMPO_URL} auth: authenticator: basicauth/grafana_cloud_tempo prometheusremotewrite/grafana_cloud_metrics: endpoint: ${env:PROMETHEUS_URL} add_metric_suffixes: false auth: authenticator: basicauth/grafana_cloud_prometheus loki/grafana_cloud_logs: endpoint: ${env:LOKI_URL} auth: authenticator: basicauth/grafana_cloud_lokidebugexporter 以 detailed 级别把收到的数据打印到 Collector 日志中便于在配置联调阶段确认数据是否真正到达otlp/grafana_cloud_traces用 Tempo 的 OTLP 上报地址与 basicauth 鉴权转发 traceprometheusremotewrite/grafana_cloud_metrics将指标以 Prometheus remote write 协议写入云端 Prometheusadd_metric_suffixes: false用于保留原始指标名避免额外追加单位后缀loki/grafana_cloud_logs负责把日志推送到云 Loki。3.5 service三条数据管线的组装service: extensions: - basicauth/grafana_cloud_tempo - basicauth/grafana_cloud_prometheus - basicauth/grafana_cloud_loki pipelines: traces: receivers: [otlp] processors: [batch] exporters: [debug, otlp/grafana_cloud_traces] metrics: receivers: [otlp] processors: [transform/add_resource_attributes_as_metric_attributes, batch] exporters: [debug, prometheusremotewrite/grafana_cloud_metrics] logs: receivers: [otlp] processors: [resource/loki, batch] exporters: [debug, loki/grafana_cloud_logs]Collector 内部按数据类型拆成三条互不干扰的 pipelinetraces 先 batch 再同时发往 debug 与 Tempometrics 先把资源属性展开为指标维度transform/add_resource_attributes_as_metric_attributes由镜像内置的 transform 处理器提供batch 后发往 debug 与 Prometheus remote writelogs 先由 resource/loki 打上格式标记再 batch 后发往 debug 与 Loki。每一条 pipeline 都保留了 debug 出口方便故障排查。3.6 docker-compose 编排与凭据注入docker-compose.yaml 使用官方 contrib 镜像内置 transform、prometheusremotewrite、loki 等社区组件将本地的 config.yaml 挂载进容器作为启动配置并暴露 4317 端口services: otel-collector: image: otel/opentelemetry-collector-contrib volumes: - ./config.yaml:/etc/otelcol-contrib/config.yaml ports: - ${OTLP_GRPC_PORT:-4317}:4317 # OTLP gRPC receiver environment: TEMPO_URL: tempo-prod-10-prod-eu-west-2.grafana.net:443 TEMPO_USERNAME: TEMPO_PASSWORD: PROMETHEUS_URL: https://prometheus-prod-24-prod-eu-west-2.grafana.net/api/prom/push PROMETHEUS_USERNAME: PROMETHEUS_PASSWORD: LOKI_URL: https://logs-prod-012.grafana.net/loki/api/v1/push LOKI_USERNAME: LOKI_PASSWORD:需要关注的环境变量如下环境变量含义取值来源OTLP_GRPC_PORT宿主机映射端口默认 4317${OTLP_GRPC_PORT:-4317}语法表示未设置时回退到 4317启动命令传入TEMPO_URLGrafana Cloud Tempo 的 OTLP 上报地址Grafana Cloud 账号的 Traces 实例详情TEMPO_USERNAME/TEMPO_PASSWORDTempo 实例 ID 与 API 令牌同上PROMETHEUS_URLPrometheus remote write 推送地址Metrics 实例的 Prometheus endpointPROMETHEUS_USERNAME/PROMETHEUS_PASSWORD指标实例 ID 与令牌同上LOKI_URLLoki push API 地址Logs 实例的 Loki endpointLOKI_USERNAME/LOKI_PASSWORD日志实例 ID 与令牌同上配置步骤登录 Grafana Cloud 后会被引导到My Account页面在那里可以为对应 Grafana 服务创建与管理各类令牌将用户名通常是实例 ID与密码/令牌分别填入 docker-compose.yaml 中 TEMPO、PROMETHEUS、LOKI 三组环境变量再将 URL 替换为实际分配给你的实例地址上面文件中的 grafana.net 主机名与路径是仓库自带的示例值应以你自己的实例信息为准。注意.yaml中的 URL 一旦包含$、:等特殊字符需遵循 YAML 引号规则或优先通过 shell 环境注入。四、启动 OpenTelemetry Collector配置完成后在 docker-compose.yaml 所在目录执行OTLP_GRPC_PORTOTLP_GRPC_PORT docker-compose up其中OTLP_GRPC_PORT替换为你希望宿主机对外开放的端口。如果直接沿用默认端口 4317可以省略该环境变量直接运行docker-compose updocker-compose up启动后 Collector 会在 4317 端口等待来自 Pathway 管线的 OTLP/gRPC 数据。可以用docker-compose logs -f观察 Collector 日志中 debug exporter 打印的详细遥测内容以验证数据链路是否打通。五、在 Pathway 侧开启监控并运行示例管线5.1 开启监控的 API当 Collector 就绪后在管线代码中开启监控。仓库示例 monitoring_demo.py 给出最小化的两行配置pw.set_license_key(keyYOUR-KEY) pw.set_monitoring_config(server_endpointhttp://localhost:OTLP_GRPC_PORT)结合 python/pathway/internals/config.py 的实现可以了解两点关键信息set_monitoring_config仅接受关键字参数其中server_endpoint是OTLP 兼容且支持 gRPC 协议的服务器端点传None可清除已有配置它还有一个可选参数detailed_metrics_dir可把详细指标以 SQLite 数据库形式导出到指定目录该调用最终写入的是全局PathwayConfig中的monitoring_server字段在纯环境变量驱动的场景下也可以通过设置PATHWAY_MONITORING_SERVER环境变量达到同样效果见 config.py 的解析逻辑。因此当你运行docker-compose up且端口为 4317 时管线中应写成pw.set_monitoring_config(server_endpointhttp://localhost:4317)并把set_license_key的参数替换为你在 Pathway 官网申请的真实 Key。5.2 示例管线工作原理配套的 monitoring_demo.py 是一个刻意保持简单、方便观察监控效果的实时管线import logging import time import pathway as pw logging.basicConfig(levellogging.INFO) pw.set_license_key(keyYOUR-KEY) pw.set_monitoring_config(server_endpointhttp://localhost:OTLP_GRPC_PORT) class DemoStream(pw.io.python.ConnectorSubject): def run(self): while True: logging.info(Producing value) self.next(value1) time.sleep(1) class InputSchema(pw.Schema): value: int table pw.io.python.read(DemoStream(), schemaInputSchema) table table.reduce(sumpw.reducers.sum(pw.this.value)) pw.io.null.write(table) pw.run()该管线的执行逻辑如下定义DemoStream继承pw.io.python.ConnectorSubject在其run方法中每秒打印一条日志并向流中推入数值1模拟持续产生输入的实时数据源声明InputSchema单字段value: int作为连接器输出与表结构的契约用pw.io.python.read读取该流经reduce配合pw.reducers.sum对value做持续累加通过pw.io.null.write丢弃结果仅用于让计算图真实执行最后pw.run()启动无限运行的实时引擎。修改其中的 license key 与监控端点后直接执行python monitoring_demo.py运行期间每秒产生的日志既会出现在控制台示例中配置了logging.basicConfig(levellogging.INFO)也会经由 Collector 转发到 Grafana Cloud Loki管线的运行指标与每轮输入的 trace 则会分别进入云端 Prometheus 与 Tempo。六、导入 Grafana Dashboard 并查看监控结果6.1 导入步骤仓库提供了可直接导入的可视化模板 grafana-dashboard.json导入步骤与模板文档一致在浏览器中打开 Grafana Cloud进入Dashboards区域点击Import上传grafana-dashboard.json文件点击Import加载示例仪表盘。导入时如果 Grafana 提示选择数据源请分别把面板关联到你在第一步配置好的 Prometheus、Loki 与 Tempo 数据源。6.2 仪表盘包含的面板从 dashboard JSON 的内容看这份仪表盘标题为 Pathway Monitoring通过三种数据源组织面板指标类数据源为 Prometheus以avg by(root_trace_id) (...)聚合查询为核心覆盖Input Latency输入延迟、Output Latency输出延迟、Process CPU utime进程 CPU 用户态时间、Memory usage内存占用四个核心指标对应的底层指标名分别为latency_input、latency_output、process_cpu_utime、process_memory_usage且都按root_trace_id分组——这意味着你可以把一条端到端数据处理的各阶段指标对齐到同一条根 trace 上观察日志类数据源为 Loki一个 Logs 面板按service_name等标签检索并展示管线日志链路类数据源为 Tempo一个 Traces 面板用于浏览与检索最近一段时间内管线的 trace。在管线运行一段时间后回到仪表盘刷新即可看到首个批次的资源使用数据CPU、内存以及输入/输出延迟随时间变化的曲线日志与 trace 面板也会开始累积内容。下图即该监控示例所展示的可视化效果示意七、小结与排障提示本方案将 Pathway 管线的日志、指标与 trace 三路遥测统一收敛到 OpenTelemetry Collector再分流写入 Grafana Cloud 的 Loki、Prometheus 与 Tempo最后用仓库自带的 dashboard 完成可视化形成一套采集—转发—存储—展示闭环的实时监控方案。写作与排障时以下几个要点值得牢记Pathway 侧set_monitoring_config只认OTLP/gRPC端点端口默认 4317若 Collector 改用了其他宿主端口管线代码中的server_endpoint必须同步修改且该能力依赖有效的 Pathway Scale 级 License KeyCollector 的三条 pipeline 都挂了debugexporter联调阶段可用docker-compose logs直接确认日志、指标与 trace 是否成功进入 Collector所有云端凭据通过环境变量注入 docker-compose.yaml切勿把真实令牌写入 config.yaml若希望进一步研究指标语义与导出细节可继续阅读 python/pathway/internals/graph_runner/telemetry.py遥测数据的组装与 OTLP exporter 初始化以及 python/pathway/internals/config.py配置项与对应环境变量将监控项与源码中的服务资源属性service.name、run.id、license.key等对应起来理解。【免费下载链接】pathwayPython ETL framework for stream processing, real-time analytics, LLM pipelines, and RAG.项目地址: https://gitcode.com/GitHub_Trending/pa/pathway创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考