可观测性最容易变成一张工具清单:Prometheus 收指标,Loki 收日志,Jaeger 收 Trace,Grafana 做大盘。容器都能启动,不代表生产问题就能被定位。
我重新对照了 GoChat 的代码和部署配置。结论是:项目已经有结构化日志、Prometheus/Loki/Grafana 容器和一个 OpenTelemetry 封装,但“三种信号串起来”还是设计目标,不是已经完成的现状。
这篇文章先记录审计结果,再给出我会如何补齐它。
先定义要回答的问题
IM 系统出故障时,用户的描述通常只有“消息发不出去”、“对方很久才收到”或“一直重连”。可观测性体系需要将这些现象逐层分解:
- 影响了多少用户,从什么时候开始?
- 问题在握手、Gateway 发送队列、Kafka、Logic、Task,还是下行 Gateway?
- 是所有消息都变慢,还是某个 topic、节点、会话类型或版本?
- 错误是否持续,重试和重连有没有进一步放大流量?
- 单条问题消息经历了哪些步骤,每一段花了多久?
Metrics 适合回答范围和趋势,Trace 适合回答一条请求或消息的路径,Logs 适合保留离散事件和业务上下文。三者需要共享 service.name、service.instance.id、环境、版本和 Trace ID,否则仍然是三座孤岛。
当前仓库实际有什么
日志:已有结构化输出,采集链路需要迁移
GoChat 的 im-infra/clog 支持 namespace、JSON 输出和从 Context 提取 trace_id。Gateway、Logic、Task 和 Repo 已经在多个模块中使用 clog.Namespace。
统一 Docker Compose 通过 Promtail 读取 Docker JSON log,写入 Loki,Grafana 已配置 Loki 数据源。这条链路的基本方向是正确的,但有两个现实问题:
- Promtail 已于 2026 年 3 月 2 日结束生命周期,Grafana 建议迁移到 Alloy 或其他受支持的客户端。
- 当前 Promtail 只按 Docker 日志文件采集,没有解析 JSON 并提升
service、level等低基数字段,查询时会过度依赖全文搜索。
迁移时不要把 user_id、message_id 和 trace_id 都变成 Loki label。这些字段基数很高,应保留在 JSON 内容中查询;label 保留 service、environment、level 和 instance 等可控维度。
Metrics:有封装代码,但统一 Prometheus 没有抓取业务服务
im-infra/metrics 封装了 OpenTelemetry TracerProvider、MeterProvider、HTTP/gRPC 拦截器、Counter 和 Histogram。这部分代码是一个可用起点,但当前业务服务没有 metrics.New(...) 初始化调用。
统一 Prometheus 配置也只抓取 localhost:9090,也就是 Prometheus 自己:
scrape_configs:
- job_name: prometheus
static_configs:
- targets: [localhost:9090]
仓库虽然有 CPU、内存、HTTP 错误率和延迟告警规则,Prometheus 主配置没有通过 rule_files 加载它们;同时也没有 node_exporter 抓取任务,因此 node_cpu_seconds_total 等表达式没有数据源。
“配置文件存在”和“告警真的会触发”之间,还差了指标注册、端点暴露、抓取配置、规则加载和 Alertmanager 路由。
Tracing:有 Provider 和手动 Trace ID,还没有统一链路
im-infra/metrics 可以导出到 Jaeger、Zipkin 或 stdout,也会设置 W3C Trace Context propagator。但统一部署配置没有 Jaeger、Tempo 或 OpenTelemetry Collector,业务服务也没有初始化 Provider。
Kafka 消息结构已经有 TraceID,生产者和消费者的部分代码会手动传播该值。这能用来关联日志,但它不等同于完整的 OpenTelemetry Context 传播:Span ID、trace flags 和 baggage 并没有一起传递,消费端也没有因此创建正确的 consumer span/link。
把“已有”和“目标”放到一张表里
| 信号 | 当前状态 | 下一步 |
|---|---|---|
| Logs | clog JSON/namespace,Promtail -> Loki -> Grafana | 迁移 Alloy,统一 service/version/instance 字段 |
| Metrics | OTel 封装代码,零散的 /metrics 端点 | 在每个服务初始化,补 scrape jobs 和真实业务指标 |
| Traces | Provider 代码与手动 TraceID 字段 | 使用 OTLP -> Collector -> Trace backend,打通 Kafka 传播 |
| Alerts | 有未加载的 rules 文件 | 加载规则,补 Alertmanager,用故障注入验证 |
| Correlation | 部分日志包含 TraceID | Grafana 中实现 metrics -> traces -> logs 跳转 |
目标拓扑:让服务只对接 OTLP
新的部署不再让每个 Go 服务分别知道 Jaeger、Zipkin 和后端凭据。OpenTelemetry 官方在生产环境通常建议将 Collector 放在服务与后端之间,让 Collector 处理 batch、retry、过滤和敏感字段。
Go services
-> OTLP
-> OpenTelemetry Collector
-> trace backend (Tempo or Jaeger)
-> metrics backend / Prometheus integration
Docker logs
-> Grafana Alloy
-> Loki
Prometheus + Loki + trace backend
-> Grafana
小规模开发环境可以直接导出到后端,先验证埋点。生产环境引入 Collector 后,应用端只保留 OTLP endpoint、TLS 和资源字段配置。
第一步:统一服务身份
每条 metric、log 和 span 至少应该共享:
service.name
service.version
service.instance.id
deployment.environment.name
Kubernetes 中的 pod.name、namespace、node.name 可以由 Collector/Alloy 的资源检测与元数据处理器补充。应用代码不应把 Pod 名写死在配置中。
这一层如果不统一,Grafana 里的 service="im-gateway"、service_name="gateway" 和 job="gochat-gateway" 会指向同一个服务,Dashboard 和跳转都会变得脆弱。
第二步:先打通 HTTP/gRPC,再处理 Kafka
HTTP 和 gRPC 已经有成熟的 OpenTelemetry instrumentation。统一 bootstrap 应该在每个服务启动时初始化 Provider,在停机时用带 deadline 的 Context flush。
先验证一条同步链路:
client -> gateway HTTP -> logic gRPC -> repo gRPC/database
当 Trace 后端能看到正确的父子 span、status 和 service name 后,再接 Kafka。这样能把“Provider 没启动”、“Context 没传递”和“异步语义不对”分开排查。
第三步:用 propagator 传播 Kafka Context
上行消息已有 Headers map[string]string,可以使用 OpenTelemetry propagator 注入 W3C Trace Context:
headers := make(map[string]string)
otel.GetTextMapPropagator().Inject(
ctx,
propagation.MapCarrier(headers),
)
消费端先提取 Context,再创建 consumer span:
ctx = otel.GetTextMapPropagator().Extract(
ctx,
propagation.MapCarrier(headers),
)
ctx, span := tracer.Start(ctx, "consume gochat.upstream")
defer span.End()
异步消息不总是简单父子关系。一条消息被批处理、延迟消费或合并多个输入时,可能更适合 span link。具体语义应和消费确认、重试和死信协议一起设计。
第四步:埋点围绕消息生命周期
IM 系统需要的不只是 Go runtime 和 HTTP 默认指标。我会先补下面这组业务指标:
Gateway 连接
gochat_gateway_connections{instance}
gochat_gateway_handshakes_total{result,reason}
gochat_gateway_disconnects_total{reason}
gochat_gateway_send_queue_full_total{instance}
gochat_gateway_reconnects_total{client_version}
connections 是 Gauge,其他多数是 Counter。user_id、IP 和 connection ID 不能作为 metric label,否则时序数量会随用户数增长。
消息链路
gochat_messages_total{stage,result,message_type}
gochat_message_stage_duration_seconds{stage}
gochat_message_e2e_duration_seconds{conversation_type}
gochat_message_retries_total{stage,reason}
gochat_message_duplicates_total{stage}
stage 可以是 gateway_ingress、logic、task_fanout、gateway_egress。端到端延迟要统一起止时刻,并注意节点时钟偏差。如果用客户端时间做起点,需要明确它不适合作为服务端精确 SLA 根据。
Kafka
gochat_kafka_produce_total{topic,result}
gochat_kafka_produce_duration_seconds{topic}
gochat_kafka_consumer_lag{topic,partition,group}
gochat_kafka_consume_total{topic,result}
Topic、partition 和 consumer group 通常是可控维度,但动态为每个 Gateway 创建 topic/group 时要计算时序基数。用户路由不应直接变成指标 label。
告警从用户体验开始
现有告警规则以 CPU、内存和 HTTP 5xx 为主。这些运维指标很有用,但它们不能直接表示用户是否收到消息。
第一批告警可以围绕:
- WebSocket 握手失败率和异常断开率。
- 消息入站成功率、下行投递成功率和端到端 P99 延迟。
- Kafka consumer lag 的持续增长,不只是某一时刻超过固定值。
- Gateway 发送队列满、重连风暴和单节点连接数偏斜。
例如消息错误率:
sum(rate(gochat_messages_total{result="error"}[5m]))
/
clamp_min(sum(rate(gochat_messages_total[5m])), 1)
> 0.05
P99 延迟:
histogram_quantile(
0.99,
sum by (le) (
rate(gochat_message_e2e_duration_seconds_bucket[5m])
)
)
> 1
阈值只是示例。真实阈值应来自服务目标、基线数据和告警响应能力。如果每天正常流量峰值都会报警,团队很快就会学会忽略它。
日志只记录有新信息的事件
一条业务消息经过多个服务,如果每一层都打“开始处理”和“处理完成”,高吞吐下日志量会很快失控。
建议保留:
- 连接建立、鉴权失败、被心跳超时和队列满关闭。
- 消息状态转换、重试、死信和幂等冲突。
- 外部依赖错误、超时、熔断和结构化错误码。
- 慢路径的摘要,详细时间线交给 Trace。
消息正文、token、密码、完整 SQL 和请求体不应默认写入日志。脱敏和丢弃超大字段最好放在 clog Handler 或 Collector/Alloy processor 中统一完成。
采样不能只写一个 1%
头部采样在请求开始时做决定,成本可预测,但当时还不知道请求最后是否慢或失败。尾部采样由 Collector 在收集到一条 Trace 的 span 后决定,可以优先保留错误、高延迟和特定业务路径,但需要为 Collector 的内存、等待时间和水平扩展付出成本。
在 GoChat 尚未打通 Trace 时,先在开发/压测环境用 100% 采样检查语义,然后根据真实 QPS 和 span 数量设计生产策略。不要在还看不到一条完整链路时,先将数据随机丢掉 99%。
一条可执行的落地顺序
- 统一 service/version/instance/environment 资源字段。
- 修正 Prometheus
scrape_configs和rule_files,确保一个真实应用 Counter 能出现在 Grafana。 - 在 Gateway -> Logic -> Repo 同步路径初始化 OTel,通过 Collector 导出 Trace。
- 用 W3C Trace Context 打通 Kafka 生产和消费,确认重试与 batch 下的 span/link 语义。
- 补 Gateway 连接、消息各阶段、Kafka lag 和端到端延迟指标。
- 将 Promtail 迁移到 Alloy,保留低基数 label,并在日志中注入 Trace ID。
- 在 Grafana 配置 metrics -> traces -> logs 跳转和最小可用 Dashboard。
- 用停 Kafka consumer、增加网络延迟、填满 Gateway 队列等故障注入验证告警和排查手册。
真正的完成标准可以用一次故障排查来检验:消息延迟告警触发后,值班人员能进入异常时间段,找到一条慢 Trace,跳到同一 Trace ID 的日志,最后确认出问题的节点或 topic。