From df96ec88792f70326a74fe9a21914963baa3e4b2 Mon Sep 17 00:00:00 2001 From: lingniu Date: Thu, 2 Jul 2026 19:27:36 +0800 Subject: [PATCH] feat(go): expose nats bridge pending metrics --- docs/ops/go-service-observability.md | 3 ++ .../cmd/nats-kafka-bridge/main.go | 30 +++++++++++++++++-- .../cmd/nats-kafka-bridge/main_test.go | 23 ++++++++++++++ 3 files changed, 54 insertions(+), 2 deletions(-) diff --git a/docs/ops/go-service-observability.md b/docs/ops/go-service-observability.md index 5cbcd33f..07c3ae8e 100644 --- a/docs/ops/go-service-observability.md +++ b/docs/ops/go-service-observability.md @@ -37,6 +37,9 @@ curl -fsS http://127.0.0.1:20200/readyz | `vehicle_bridge_messages_total` | NATS messages fetched by the bridge. Labels: `subject`, `status`. | | `vehicle_bridge_kafka_writes_total` | Bridge writes to Kafka. Labels: `topic`, `status`. | | `vehicle_bridge_nats_acks_total` | NATS ack results after Kafka write. Labels: `subject`, `status`. | +| `vehicle_bridge_nats_consumer_pending` | JetStream messages pending for the bridge durable consumer. | +| `vehicle_bridge_nats_consumer_ack_pending` | JetStream messages delivered to bridge but not yet acked. | +| `vehicle_bridge_nats_consumer_waiting` | Pull requests waiting on the bridge durable consumer. | | `vehicle_history_writes_total` | TDengine history writes. Labels: `topic`, `status`. | | `vehicle_stat_writes_total` | MySQL metric writes. Labels: `topic`, `status`. | | `vehicle_realtime_updates_total` | Redis/MySQL realtime projector updates. Labels: `topic`, `status`. | diff --git a/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go b/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go index 98e52583..a260705c 100644 --- a/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go +++ b/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go @@ -80,7 +80,7 @@ func main() { "durable", cfg.NATSDurable, "filter", cfg.NATSFilter, "kafka_brokers", strings.Join(cfg.KafkaBrokers, ",")) - runBridge(ctx, logger, registry, sub, writer, cfg) + runBridge(ctx, logger, registry, js, sub, writer, cfg) } type config struct { @@ -127,6 +127,10 @@ type natsPullSubscription interface { Fetch(int, ...nats.PullOpt) ([]*nats.Msg, error) } +type natsConsumerInfoReader interface { + ConsumerInfo(stream string, name string, opts ...nats.JSOpt) (*nats.ConsumerInfo, error) +} + type kafkaBatchWriter interface { WriteMessages(context.Context, ...kafka.Message) error } @@ -137,11 +141,23 @@ type bridgeMessage struct { ack func() error } -func runBridge(ctx context.Context, logger *slog.Logger, registry *metrics.Registry, sub natsPullSubscription, writer kafkaBatchWriter, cfg config) { +func runBridge(ctx context.Context, logger *slog.Logger, registry *metrics.Registry, infoReader natsConsumerInfoReader, sub natsPullSubscription, writer kafkaBatchWriter, cfg config) { + var lastConsumerInfoAt time.Time for { if ctx.Err() != nil { return } + if time.Since(lastConsumerInfoAt) >= time.Duration(envInt("NATS_CONSUMER_METRICS_INTERVAL_SECONDS", 10))*time.Second { + lastConsumerInfoAt = time.Now() + if infoReader != nil { + info, err := infoReader.ConsumerInfo(cfg.NATSStream, cfg.NATSDurable, nats.Context(ctx)) + if err != nil { + logger.Warn("nats consumer info failed", "stream", cfg.NATSStream, "durable", cfg.NATSDurable, "error", err) + } else { + recordNATSConsumerInfoMetrics(registry, cfg, info) + } + } + } msgs, err := sub.Fetch(cfg.BatchSize, nats.MaxWait(cfg.FetchWait)) if err != nil { if errors.Is(err, nats.ErrTimeout) { @@ -246,6 +262,16 @@ func addBridgeTopicMetric(registry *metrics.Registry, name string, topic string, registry.IncCounter(name, metrics.Labels{"topic": topic, "status": status}) } +func recordNATSConsumerInfoMetrics(registry *metrics.Registry, cfg config, info *nats.ConsumerInfo) { + if registry == nil || info == nil { + return + } + labels := metrics.Labels{"stream": cfg.NATSStream, "consumer": cfg.NATSDurable} + registry.SetGauge("vehicle_bridge_nats_consumer_pending", labels, float64(info.NumPending)) + registry.SetGauge("vehicle_bridge_nats_consumer_ack_pending", labels, float64(info.NumAckPending)) + registry.SetGauge("vehicle_bridge_nats_consumer_waiting", labels, float64(info.NumWaiting)) +} + func kafkaMessage(topic string, data []byte) kafka.Message { var env envelope.FrameEnvelope message := kafka.Message{Topic: topic, Value: data} diff --git a/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go b/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go index f4c091a2..15ba57d5 100644 --- a/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go +++ b/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go @@ -7,6 +7,7 @@ import ( "strings" "testing" + "github.com/nats-io/nats.go" "github.com/segmentio/kafka-go" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" @@ -90,6 +91,28 @@ func TestBridgeBatchRecordsMetrics(t *testing.T) { } } +func TestRecordNATSConsumerInfoMetrics(t *testing.T) { + registry := metrics.NewRegistry() + cfg := config{NATSStream: "VEHICLE_INGEST", NATSDurable: "vehicle-kafka-bridge"} + + recordNATSConsumerInfoMetrics(registry, cfg, &nats.ConsumerInfo{ + NumPending: 17, + NumAckPending: 3, + NumWaiting: 2, + }) + + text := registry.Render() + for _, want := range []string{ + `vehicle_bridge_nats_consumer_pending{consumer="vehicle-kafka-bridge",stream="VEHICLE_INGEST"} 17`, + `vehicle_bridge_nats_consumer_ack_pending{consumer="vehicle-kafka-bridge",stream="VEHICLE_INGEST"} 3`, + `vehicle_bridge_nats_consumer_waiting{consumer="vehicle-kafka-bridge",stream="VEHICLE_INGEST"} 2`, + } { + if !strings.Contains(text, want) { + t.Fatalf("metrics missing %s:\n%s", want, text) + } + } +} + type recordingBridgeWriter struct { messages []kafka.Message err error