From 3b694671fa4a332562c7cf943ece9c71c47db2fe Mon Sep 17 00:00:00 2001 From: lingniu Date: Fri, 3 Jul 2026 19:17:02 +0800 Subject: [PATCH] feat(go): batch nats fast writer tdengine writes --- docs/ops/100k-capacity-baseline.md | 2 +- docs/ops/go-vehicle-ingest-memory.md | 2 + .../cmd/nats-fast-writer/main.go | 94 ++++++++++++++++--- .../cmd/nats-fast-writer/main_test.go | 45 ++++++++- 4 files changed, 127 insertions(+), 16 deletions(-) diff --git a/docs/ops/100k-capacity-baseline.md b/docs/ops/100k-capacity-baseline.md index 22fb2afd..4f9e028d 100644 --- a/docs/ops/100k-capacity-baseline.md +++ b/docs/ops/100k-capacity-baseline.md @@ -39,7 +39,7 @@ - Gateway 默认 `TCP_MAX_CONNECTIONS` 已调整为 `120000`,生产仍可通过环境变量覆盖。 - 已新增可重复的 TCP 连接压测工具,后续需要用它跑 10K/50K/100K 阶段测试并记录结果。 - 仍缺读超时按协议维度计数;Gateway frame duration 已提供 histogram,可用于 parse + enqueue + response 的 p95/p99 估算。 -- TDengine history writer 当前逐消息插入,后续需要批量写入。 +- TDengine history writer 和 NATS fast writer 已使用 batch 写入;后续需要用带帧率压测验证 batch size、flush latency 和存储端承载能力。 - Kafka topic 当前 12 分区,100K 目标下需要结合实际 FPS 再评估分区数。 ## ECS OS 参数建议 diff --git a/docs/ops/go-vehicle-ingest-memory.md b/docs/ops/go-vehicle-ingest-memory.md index 2d5f00ae..06aac1e6 100644 --- a/docs/ops/go-vehicle-ingest-memory.md +++ b/docs/ops/go-vehicle-ingest-memory.md @@ -421,6 +421,8 @@ vehicle_fast_writer_stage_duration_ms_histogram_sum{subject,stage,status} `stage` 取值包括 `tdengine`、`redis`、`ack`,分别对应快速写 TDengine raw/location、Redis realtime 投影、NATS ack。带帧率压测时如果某个 stage 的 p99 持续上升,可以直接定位 fast path 是存储瓶颈还是 ack/网络瓶颈。 +2026-07-03 后续优化:`nats-fast-writer` 已从“Fetch batch 后逐条写 TDengine”改为“Fetch batch 后先调用 `history.Writer.AppendAllBatch` 批量写 TDengine,再逐条更新 Redis 和 ack”。这样可以减少 TDengine round trip;Redis/ack 仍按单条执行,避免某条实时投影失败时误 ack。 + ### NATS Kafka bridge batch 指标 2026-07-03 已为 NATS -> Kafka bridge 增加 batch pending 和 duration histogram: diff --git a/go/vehicle-gateway/cmd/nats-fast-writer/main.go b/go/vehicle-gateway/cmd/nats-fast-writer/main.go index 30433987..26dbd3a1 100644 --- a/go/vehicle-gateway/cmd/nats-fast-writer/main.go +++ b/go/vehicle-gateway/cmd/nats-fast-writer/main.go @@ -184,6 +184,7 @@ func loadConfig() config { type fastAppender interface { AppendAll(context.Context, envelope.FrameEnvelope) error + AppendAllBatch(context.Context, []envelope.FrameEnvelope) error } type fastUpdater interface { @@ -210,23 +211,27 @@ func runFastWorker(ctx context.Context, logger *slog.Logger, registry *metrics.R time.Sleep(time.Second) continue } + fastMessages := make([]*fastMessage, 0, len(msgs)) for _, msg := range msgs { - msg := msg - operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait) - err := processFastMessage(operationCtx, registry, appender, updater, &fastMessage{ - subject: msg.Subject, - data: msg.Data, + natsMsg := msg + fastMessages = append(fastMessages, &fastMessage{ + subject: natsMsg.Subject, + data: natsMsg.Data, ack: func() error { - return msg.Ack() + return natsMsg.Ack() }, }) - cancel() - if err != nil { - addFastMetric(registry, msg.Subject, "error") - logger.Error("fast write failed", "subject", msg.Subject, "error", err) - continue - } - addFastMetric(registry, msg.Subject, "ok") + } + operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait) + err = processFastBatch(operationCtx, registry, appender, updater, fastMessages) + cancel() + if err != nil { + addFastMetric(registry, fastBatchSubject(fastMessages), "error") + logger.Error("fast write batch failed", "messages", len(fastMessages), "error", err) + continue + } + for _, msg := range fastMessages { + addFastMetric(registry, msg.subject, "ok") } } } @@ -237,6 +242,53 @@ type natsPullSubscription interface { var fastWriterStageDurationBucketsMS = []float64{1, 5, 10, 25, 50, 100, 250, 500, 1000, 5000} +func processFastBatch(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, messages []*fastMessage) error { + if len(messages) == 0 { + return nil + } + envelopes := make([]envelope.FrameEnvelope, 0, len(messages)) + validMessages := make([]*fastMessage, 0, len(messages)) + for _, msg := range messages { + var env envelope.FrameEnvelope + if err := json.Unmarshal(msg.data, &env); err != nil { + if msg.ack != nil { + _ = msg.ack() + } + continue + } + envelopes = append(envelopes, env) + validMessages = append(validMessages, msg) + } + if len(envelopes) == 0 { + return nil + } + subject := fastBatchSubject(validMessages) + started := time.Now() + err := appender.AppendAllBatch(ctx, envelopes) + recordFastWriterStageDuration(registry, subject, "tdengine", statusFromError(err), time.Since(started)) + if err != nil { + return fmt.Errorf("tdengine batch append: %w", err) + } + for i, env := range envelopes { + msg := validMessages[i] + started = time.Now() + err = updater.FastUpdate(ctx, env) + recordFastWriterStageDuration(registry, msg.subject, "redis", statusFromError(err), time.Since(started)) + if err != nil { + return fmt.Errorf("redis fast update: %w", err) + } + if msg.ack != nil { + started = time.Now() + err = msg.ack() + recordFastWriterStageDuration(registry, msg.subject, "ack", statusFromError(err), time.Since(started)) + if err != nil { + return fmt.Errorf("nats ack: %w", err) + } + } + } + return nil +} + func processFastMessage(ctx context.Context, registry *metrics.Registry, appender fastAppender, updater fastUpdater, msg *fastMessage) error { var env envelope.FrameEnvelope if err := json.Unmarshal(msg.data, &env); err != nil { @@ -268,6 +320,22 @@ func processFastMessage(ctx context.Context, registry *metrics.Registry, appende return nil } +func fastBatchSubject(messages []*fastMessage) string { + if len(messages) == 0 { + return "unknown" + } + subject := messages[0].subject + for _, msg := range messages[1:] { + if msg.subject != subject { + return "mixed" + } + } + if strings.TrimSpace(subject) == "" { + return "unknown" + } + return subject +} + func ensureStream(js nats.JetStreamContext, cfg config) error { stream := &nats.StreamConfig{ Name: cfg.NATSStream, diff --git a/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go b/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go index e083e0cb..591fa86b 100644 --- a/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go +++ b/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go @@ -80,9 +80,44 @@ func TestProcessFastMessageRecordsStageDurationMetrics(t *testing.T) { } } +func TestProcessFastBatchAppendsTDengineBatchBeforeRedisAndAck(t *testing.T) { + first := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, VIN: "VIN001", EventID: "evt-4"} + second := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, VIN: "VIN002", EventID: "evt-5"} + firstPayload, err := first.MarshalJSONBytes() + if err != nil { + t.Fatal(err) + } + secondPayload, err := second.MarshalJSONBytes() + if err != nil { + t.Fatal(err) + } + appender := &recordingFastAppender{} + updater := &recordingFastUpdater{} + ackCount := 0 + msgs := []*fastMessage{ + {subject: "vehicle.raw.go.jt808.v1", data: firstPayload, ack: func() error { ackCount++; return nil }}, + {subject: "vehicle.raw.go.jt808.v1", data: secondPayload, ack: func() error { ackCount++; return nil }}, + } + + if err := processFastBatch(context.Background(), nil, appender, updater, msgs); err != nil { + t.Fatalf("processFastBatch() error = %v", err) + } + if appender.count != 0 { + t.Fatalf("AppendAll count = %d, want 0", appender.count) + } + if appender.batchCount != 1 || appender.batchRows != 2 { + t.Fatalf("AppendAllBatch count=%d rows=%d, want count=1 rows=2", appender.batchCount, appender.batchRows) + } + if updater.count != 2 || ackCount != 2 { + t.Fatalf("updates=%d acks=%d, want 2/2", updater.count, ackCount) + } +} + type recordingFastAppender struct { - count int - err error + count int + batchCount int + batchRows int + err error } func (a *recordingFastAppender) AppendAll(context.Context, envelope.FrameEnvelope) error { @@ -90,6 +125,12 @@ func (a *recordingFastAppender) AppendAll(context.Context, envelope.FrameEnvelop return a.err } +func (a *recordingFastAppender) AppendAllBatch(_ context.Context, envs []envelope.FrameEnvelope) error { + a.batchCount++ + a.batchRows += len(envs) + return a.err +} + type recordingFastUpdater struct { count int err error