diff --git a/docs/ops/go-vehicle-ingest-memory.md b/docs/ops/go-vehicle-ingest-memory.md index 82af8ca8..d46813e3 100644 --- a/docs/ops/go-vehicle-ingest-memory.md +++ b/docs/ops/go-vehicle-ingest-memory.md @@ -433,6 +433,15 @@ vehicle_fast_writer_stage_duration_ms_histogram_sum{subject,stage,status} 2026-07-03 后续优化:`nats-fast-writer` 已从“Fetch batch 后逐条写 TDengine”改为“Fetch batch 后先调用 `history.Writer.AppendAllBatch` 批量写 TDengine,再逐条更新 Redis 和 ack”。这样可以减少 TDengine round trip;Redis/ack 仍按单条执行,避免某条实时投影失败时误 ack。 +2026-07-03 后续优化:`nats-fast-writer` 的 TDengine SQL 连接池已从固定 `1` 改为可配置: + +```text +FAST_WRITER_TDENGINE_MAX_OPEN_CONNS +FAST_WRITER_TDENGINE_MAX_IDLE_CONNS +``` + +默认 `max_open=1`、`max_idle=1`。原因是当前 TDengine 写入依赖启动时执行的 `USE `;`database/sql` 连接池新建连接不会继承其他连接上的 `USE` 状态。实测默认跟随 worker 放大到 8 会导致 `[0x200] db is not specified` / `[0x2616] Database not specified`。后续只有在 TDengine DSN 或 writer SQL 改为每条连接都明确选库后,才能提高该连接池。 + 同时新增 batch pending gauge: ```text diff --git a/docs/ops/vehicle-ingest-runbook.md b/docs/ops/vehicle-ingest-runbook.md index 7ba03b3b..94d617d1 100644 --- a/docs/ops/vehicle-ingest-runbook.md +++ b/docs/ops/vehicle-ingest-runbook.md @@ -136,7 +136,7 @@ go run ./cmd/load-sim \ | `vehicle_fast_writer_nats_consumer_pending` | 持续增长且 `> 10000` | fast-writer 消费 NATS 的速度跟不上入口写入速度。 | | `vehicle_fast_writer_batch_pending_messages` | 持续非 0 或 burst 后不回落 | fast-writer 已拉取 NATS 消息但尚未完成 TDengine/Redis 写入和 ack。 | | `vehicle_fast_writer_batch_pending_envelopes` | 持续非 0 或 burst 后不回落 | fast-writer 当前批次已有有效 envelope 在等待落库或 ack。 | -| `vehicle_fast_writer_stage_duration_ms_histogram_bucket` | 某个 stage 的 p99 连续 5 分钟上升 | NATS 快速写链路在 TDengine、Redis 或 NATS ack 某一阶段变慢。 | +| `vehicle_fast_writer_stage_duration_ms_histogram_bucket` | 某个 stage 的 p99 连续 5 分钟上升 | NATS 快速写链路在 TDengine、Redis 或 NATS ack 某一阶段变慢;当前 TDengine 连接池默认保持 `1`,放大 `FAST_WRITER_TDENGINE_MAX_OPEN_CONNS` 前必须先确认每条连接都能明确选库。 | | `vehicle_realtime_store_update_duration_ms_histogram_bucket{store="redis"}` | p99 连续 5 分钟上升 | Redis 实时投影变慢,会直接影响 realtime consumer 追平能力。 | | `vehicle_realtime_store_update_duration_ms_histogram_bucket{store="mysql"}` | p99 连续 5 分钟上升 | MySQL 当前态/位置投影变慢,需结合 async queue depth 和 dropped 计数判断。 | | `vehicle_stat_write_duration_ms_histogram_bucket` | p99 连续 5 分钟上升 | MySQL 每日里程统计写入变慢,可能导致 stat Kafka lag 增长。 | diff --git a/go/vehicle-gateway/cmd/nats-fast-writer/main.go b/go/vehicle-gateway/cmd/nats-fast-writer/main.go index 7eb95da8..e7fc1013 100644 --- a/go/vehicle-gateway/cmd/nats-fast-writer/main.go +++ b/go/vehicle-gateway/cmd/nats-fast-writer/main.go @@ -66,8 +66,8 @@ func main() { os.Exit(1) } defer tdDB.Close() - tdDB.SetMaxOpenConns(1) - tdDB.SetMaxIdleConns(1) + tdDB.SetMaxOpenConns(cfg.TDengineMaxOpenConns) + tdDB.SetMaxIdleConns(cfg.TDengineMaxIdleConns) if err := tdDB.PingContext(ctx); err != nil { logger.Error("tdengine ping failed", "error", err) os.Exit(1) @@ -117,6 +117,8 @@ func main() { "filter", cfg.NATSFilter, "batch_size", cfg.BatchSize, "workers", cfg.Workers, + "tdengine_max_open_conns", cfg.TDengineMaxOpenConns, + "tdengine_max_idle_conns", cfg.TDengineMaxIdleConns, "operation_timeout_ms", cfg.OperationWait.Milliseconds()) for i := 0; i < cfg.Workers; i++ { go runFastWorker(ctx, logger, registry, js, sub, historyWriter, realtimeRepo, cfg) @@ -141,6 +143,8 @@ type config struct { TDengineDSN string TDengineDatabase string TDengineEnsureSchema bool + TDengineMaxOpenConns int + TDengineMaxIdleConns int RedisAddr string RedisUsername string RedisPassword string @@ -157,6 +161,9 @@ func loadConfig() config { env("NATS_SUBJECT_JT808_FIELDS", env("KAFKA_TOPIC_JT808_FIELDS", topics.FieldsJT808)), env("NATS_SUBJECT_YUTONG_MQTT_FIELDS", env("KAFKA_TOPIC_YUTONG_MQTT_FIELDS", topics.FieldsYutongMQTT)), }, ","))) + workers := envInt("FAST_WRITER_WORKERS", 8) + maxOpenConns := envInt("FAST_WRITER_TDENGINE_MAX_OPEN_CONNS", 1) + maxIdleConns := envInt("FAST_WRITER_TDENGINE_MAX_IDLE_CONNS", maxOpenConns) return config{ NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"), NATSClientName: env("NATS_CLIENT_NAME", "lingniu-nats-fast-writer"), @@ -169,11 +176,13 @@ func loadConfig() config { OperationWait: time.Duration(envInt("FAST_WRITER_OPERATION_TIMEOUT_MS", 100)) * time.Millisecond, AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 30)) * time.Second, StreamMaxAge: time.Duration(envInt("NATS_STREAM_MAX_AGE_HOURS", 24)) * time.Hour, - Workers: envInt("FAST_WRITER_WORKERS", 8), + Workers: workers, TDengineDriver: env("TDENGINE_DRIVER", "taosWS"), TDengineDSN: env("TDENGINE_DSN", ""), TDengineDatabase: env("TDENGINE_DATABASE", history.DefaultDatabase), TDengineEnsureSchema: env("TDENGINE_ENSURE_SCHEMA", "true") != "false", + TDengineMaxOpenConns: maxOpenConns, + TDengineMaxIdleConns: maxIdleConns, RedisAddr: env("REDIS_ADDR", "127.0.0.1:6379"), RedisUsername: env("REDIS_USERNAME", ""), RedisPassword: env("REDIS_PASSWORD", ""), 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 4052a9c5..2f67de67 100644 --- a/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go +++ b/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go @@ -197,6 +197,35 @@ func TestRecordFastNATSConsumerInfoMetrics(t *testing.T) { } } +func TestLoadConfigDefaultsTDenginePoolToSingleConnection(t *testing.T) { + t.Setenv("FAST_WRITER_WORKERS", "12") + t.Setenv("FAST_WRITER_TDENGINE_MAX_OPEN_CONNS", "") + + cfg := loadConfig() + + if cfg.TDengineMaxOpenConns != 1 { + t.Fatalf("TDengineMaxOpenConns = %d, want 1", cfg.TDengineMaxOpenConns) + } + if cfg.TDengineMaxIdleConns != 1 { + t.Fatalf("TDengineMaxIdleConns = %d, want 1", cfg.TDengineMaxIdleConns) + } +} + +func TestLoadConfigReadsTDenginePoolOverride(t *testing.T) { + t.Setenv("FAST_WRITER_WORKERS", "12") + t.Setenv("FAST_WRITER_TDENGINE_MAX_OPEN_CONNS", "4") + t.Setenv("FAST_WRITER_TDENGINE_MAX_IDLE_CONNS", "2") + + cfg := loadConfig() + + if cfg.TDengineMaxOpenConns != 4 { + t.Fatalf("TDengineMaxOpenConns = %d, want 4", cfg.TDengineMaxOpenConns) + } + if cfg.TDengineMaxIdleConns != 2 { + t.Fatalf("TDengineMaxIdleConns = %d, want 2", cfg.TDengineMaxIdleConns) + } +} + type recordingFastAppender struct { count int batchCount int