如何为 pgwatch 定制 gRPC Sink:对接自有时序存储的完整开发者指南
【免费下载链接】pgwatch🔬pgwatch: PostgreSQL metrics monitor/dashboard项目地址: https://gitcode.com/gh_mirrors/pg/pgwatch
pgwatch 是一款功能强大的PostgreSQL 指标监控与仪表盘开源项目,内置了 PostgreSQL、Prometheus、JSON 文件等存储后端(Sink)。当你需要把监控指标推送到 InfluxDB、ClickHouse、IoTDB 等自有时序数据库时,pgwatch 的gRPC Sink就是官方预留的"万能出口"——只需实现一个约 5 个 RPC 接口的小服务,就能让 PostgreSQL 监控数据自由落袋。本文带你从零定制属于自己的 gRPC Sink。
📌 gRPC Sink 解决什么问题?
pgwatch 的数据链路是:Reaper(收割器)定期从被监控的 PostgreSQL 实例采集指标 → 通过 Sink 写入存储 → Web UI / 仪表盘展示。官方文档在 docs/concept/components.md 中明确说明了各 Sink 的定位,其中 gRPC Sink 的说明是:
当内置存储选项都不满足需求时,你可以基于 pgwatch 的 protobuf 定义实现自己的 gRPC 服务器,将指标推送到任意自定义存储后端。
这意味着:
- ✅不受限于官方支持的数据库——任何语言、任何时序库都能接
- ✅协议是 protobuf——高性能、强类型、跨语言(Go、Java、Python、Rust 均可实现)
- ✅附带元数据能力——不仅能推数据,还能同步指标定义与监控源的增删
🧩 先看懂 protobuf 契约:只有 3 个 RPC
整个 gRPC 契约就定义在 api/pb/pgwatch.proto 中,服务名为Receiver,包含 3 个方法:
| RPC 方法 | 请求类型 | 作用 | 你应如何处理 |
|---|---|---|---|
UpdateMeasurements | MeasurementEnvelope | 推送一批指标数据(含 DB 名、指标名、自定义标签、结构化数据行) | 核心:将数据写入你的时序库 |
SyncMetric | SyncReq | 同步指标/监控源的生命周期操作(AddOp/DeleteOp/DefineOp) | 更新本地元数据:新增/删除某实例的某指标 |
DefineMetrics | google.protobuf.Struct | 推送完整的指标定义(JSON 化后的结构) | 保存指标 schema,便于建表/建模板 |
三个方法的返回值统一是Reply,其中只有一个logmsg字段——pgwatch 会把这个字段的内容打进日志,这是官方给你的"调试回音",联调时非常好用。
数据载荷的结构(MeasurementEnvelope)大致是:
DBName:被监控的 PostgreSQL 实例名MetricName:指标名(如pg_stat_statements、pg_settings)CustomTags:在 pgwatch 中为该实例打的自定义标签Data:repeated google.protobuf.Struct,即每行一条map<string, value>的测量记录
💡 提示:
internal/sinks/rpc.go中可以看到客户端侧的完整实现逻辑——指标定义是"先转 JSON、再转 Struct"发送的,你按map<string, any>接收即可,无需关心 Go 类型细节。
⚙️ 一行配置:把 pgwatch 指向你的服务器
服务写好之后,只需在 pgwatch 启动时加一个--sink参数(可以重复使用,与内置 Sink 并行写入)。URI 格式见 docs/reference/sinks_options.md:
# 最简用法 --sink=grpc://localhost:5000/ # 带认证 + TLS --sink=grpc://user:pwd@localhost:5000/?sslrootca=/home/user/ca.crt关键特性:
- 🔐认证:
user:pwd会以 gRPC metadata 形式("username"/"password"两个字段)转发给你的服务器,服务器端自行校验,不做任何强制 - 🔒TLS:指定
sslrootca参数加载 CA 证书后自动启用加密,否则走明文 - 🏓连通性探测:pgwatch 启动时会用一次
SyncMetric空请求做 Ping,只有服务端可用才认为 Sink 就绪(Unavailable状态会立即报错退出)
代码侧的入口在internal/sinks/multiwriter.go:Sink 类型解析支持grpc和rpc两种写法,都会路由到NewRPCWriter。
🏗️ 实现你的 gRPC 服务器:4 步走
第 1 步:生成客户端代码
拿 api/pb/pgwatch.proto 生成你所在语言的桩代码。以 Go 为例:
protoc --go_out=. --go-grpc_out=. api/pb/pgwatch.proto第 2 步:实现Receiver服务
以 Go 为例,核心骨架(示意):
func (s *MyServer) UpdateMeasurements(ctx context.Context, env *pb.MeasurementEnvelope) (*pb.Reply, error) { // 1. 从 ctx 的 metadata 中取出 username/password 做认证 // 2. 把 env.Data 中每行 Struct 转成 map,写入你的时序库 // 3. 返回 &pb.Reply{Logmsg: "写入 " + env.DBName} return &pb.Reply{}, nil }第 3 步:处理元数据同步
SyncMetric里根据Operation枚举维护"哪些实例在监控哪些指标";DefineMetrics收到的是指标定义的 JSON 结构,建议落盘或入库,后续建表/建保留策略都能派上用场。
第 4 步:监听端口并返回 logmsg
返回有意义的logmsg(比如"已写入 128 行到 tsdb"),它会被 pgwatch 原样打印到日志,是验证数据链路是否打通的最快方式。
📚 官方另有一个社区维护的 gRPC 服务器集合
pgwatch-contrib/rpc(见 docs/howto/implement_grpc_server.md),针对常见存储方案提供了示例实现与搭建教程,非常适合作为生产级服务的起点参考。
🧪 验证数据链路:看 Web UI 是否健康
配置生效后,打开 pgwatch Web UI 确认 Reaper 与 Sink 均工作正常:
数据落库后,你就可以在监控大盘中看到类似下面的指标网格与性能分析视图:
✅ 避坑清单与最佳实践
| 场景 | 建议 |
|---|---|
| 批量写入 | UpdateMeasurements一次携带一整批Data,务必批量插入你的时序库,避免逐行写入拖慢 Reaper |
| 错误处理 | gRPC 错误会原样回传并记录日志;写入失败请返回明确的codes状态码,方便 pgwatch 侧定位 |
| 认证 | 若走内网可暂不认证,跨网段/公网部署务必启用 TLS(sslrootca)并在服务端校验 metadata 中的凭据 |
| 标签利用 | CustomTags是 pgwatch 实例维度的标签,正好可映射为你时序库的 tag 维度,做实例级筛选 |
| 保留策略 | SyncMetric的DeleteOp提示某实例的某指标被移除,是清理历史数据的天然钩子 |
| 联调排障 | 全程善用Reply.logmsg——pgwatch 日志里搜grpc关键字即可看到每批数据的行数与耗时 |
小结
pgwatch 的 gRPC Sink 把"对接自有时序存储"这件事压缩成了三件事:读懂 api/pb/pgwatch.proto 中的 3 个 RPC、实现你的Receiver服务器、用一行--sink=grpc://...完成接入。契约简单、协议标准,无论你的存储是 InfluxDB、ClickHouse 还是自研方案,一个下午就能让 PostgreSQL 监控指标流进自己的数据仓库。🎯
【免费下载链接】pgwatch🔬pgwatch: PostgreSQL metrics monitor/dashboard项目地址: https://gitcode.com/gh_mirrors/pg/pgwatch
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考