fix(go): keep fast writer tdengine pool safe by default
This commit is contained in:
@@ -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` 已从“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>`;`database/sql` 连接池新建连接不会继承其他连接上的 `USE` 状态。实测默认跟随 worker 放大到 8 会导致 `[0x200] db is not specified` / `[0x2616] Database not specified`。后续只有在 TDengine DSN 或 writer SQL 改为每条连接都明确选库后,才能提高该连接池。
|
||||||
|
|
||||||
同时新增 batch pending gauge:
|
同时新增 batch pending gauge:
|
||||||
|
|
||||||
```text
|
```text
|
||||||
|
|||||||
@@ -136,7 +136,7 @@ go run ./cmd/load-sim \
|
|||||||
| `vehicle_fast_writer_nats_consumer_pending` | 持续增长且 `> 10000` | fast-writer 消费 NATS 的速度跟不上入口写入速度。 |
|
| `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_messages` | 持续非 0 或 burst 后不回落 | fast-writer 已拉取 NATS 消息但尚未完成 TDengine/Redis 写入和 ack。 |
|
||||||
| `vehicle_fast_writer_batch_pending_envelopes` | 持续非 0 或 burst 后不回落 | fast-writer 当前批次已有有效 envelope 在等待落库或 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="redis"}` | p99 连续 5 分钟上升 | Redis 实时投影变慢,会直接影响 realtime consumer 追平能力。 |
|
||||||
| `vehicle_realtime_store_update_duration_ms_histogram_bucket{store="mysql"}` | p99 连续 5 分钟上升 | MySQL 当前态/位置投影变慢,需结合 async queue depth 和 dropped 计数判断。 |
|
| `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 增长。 |
|
| `vehicle_stat_write_duration_ms_histogram_bucket` | p99 连续 5 分钟上升 | MySQL 每日里程统计写入变慢,可能导致 stat Kafka lag 增长。 |
|
||||||
|
|||||||
@@ -66,8 +66,8 @@ func main() {
|
|||||||
os.Exit(1)
|
os.Exit(1)
|
||||||
}
|
}
|
||||||
defer tdDB.Close()
|
defer tdDB.Close()
|
||||||
tdDB.SetMaxOpenConns(1)
|
tdDB.SetMaxOpenConns(cfg.TDengineMaxOpenConns)
|
||||||
tdDB.SetMaxIdleConns(1)
|
tdDB.SetMaxIdleConns(cfg.TDengineMaxIdleConns)
|
||||||
if err := tdDB.PingContext(ctx); err != nil {
|
if err := tdDB.PingContext(ctx); err != nil {
|
||||||
logger.Error("tdengine ping failed", "error", err)
|
logger.Error("tdengine ping failed", "error", err)
|
||||||
os.Exit(1)
|
os.Exit(1)
|
||||||
@@ -117,6 +117,8 @@ func main() {
|
|||||||
"filter", cfg.NATSFilter,
|
"filter", cfg.NATSFilter,
|
||||||
"batch_size", cfg.BatchSize,
|
"batch_size", cfg.BatchSize,
|
||||||
"workers", cfg.Workers,
|
"workers", cfg.Workers,
|
||||||
|
"tdengine_max_open_conns", cfg.TDengineMaxOpenConns,
|
||||||
|
"tdengine_max_idle_conns", cfg.TDengineMaxIdleConns,
|
||||||
"operation_timeout_ms", cfg.OperationWait.Milliseconds())
|
"operation_timeout_ms", cfg.OperationWait.Milliseconds())
|
||||||
for i := 0; i < cfg.Workers; i++ {
|
for i := 0; i < cfg.Workers; i++ {
|
||||||
go runFastWorker(ctx, logger, registry, js, sub, historyWriter, realtimeRepo, cfg)
|
go runFastWorker(ctx, logger, registry, js, sub, historyWriter, realtimeRepo, cfg)
|
||||||
@@ -141,6 +143,8 @@ type config struct {
|
|||||||
TDengineDSN string
|
TDengineDSN string
|
||||||
TDengineDatabase string
|
TDengineDatabase string
|
||||||
TDengineEnsureSchema bool
|
TDengineEnsureSchema bool
|
||||||
|
TDengineMaxOpenConns int
|
||||||
|
TDengineMaxIdleConns int
|
||||||
RedisAddr string
|
RedisAddr string
|
||||||
RedisUsername string
|
RedisUsername string
|
||||||
RedisPassword 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_JT808_FIELDS", env("KAFKA_TOPIC_JT808_FIELDS", topics.FieldsJT808)),
|
||||||
env("NATS_SUBJECT_YUTONG_MQTT_FIELDS", env("KAFKA_TOPIC_YUTONG_MQTT_FIELDS", topics.FieldsYutongMQTT)),
|
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{
|
return config{
|
||||||
NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"),
|
NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"),
|
||||||
NATSClientName: env("NATS_CLIENT_NAME", "lingniu-nats-fast-writer"),
|
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,
|
OperationWait: time.Duration(envInt("FAST_WRITER_OPERATION_TIMEOUT_MS", 100)) * time.Millisecond,
|
||||||
AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 30)) * time.Second,
|
AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 30)) * time.Second,
|
||||||
StreamMaxAge: time.Duration(envInt("NATS_STREAM_MAX_AGE_HOURS", 24)) * time.Hour,
|
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"),
|
TDengineDriver: env("TDENGINE_DRIVER", "taosWS"),
|
||||||
TDengineDSN: env("TDENGINE_DSN", ""),
|
TDengineDSN: env("TDENGINE_DSN", ""),
|
||||||
TDengineDatabase: env("TDENGINE_DATABASE", history.DefaultDatabase),
|
TDengineDatabase: env("TDENGINE_DATABASE", history.DefaultDatabase),
|
||||||
TDengineEnsureSchema: env("TDENGINE_ENSURE_SCHEMA", "true") != "false",
|
TDengineEnsureSchema: env("TDENGINE_ENSURE_SCHEMA", "true") != "false",
|
||||||
|
TDengineMaxOpenConns: maxOpenConns,
|
||||||
|
TDengineMaxIdleConns: maxIdleConns,
|
||||||
RedisAddr: env("REDIS_ADDR", "127.0.0.1:6379"),
|
RedisAddr: env("REDIS_ADDR", "127.0.0.1:6379"),
|
||||||
RedisUsername: env("REDIS_USERNAME", ""),
|
RedisUsername: env("REDIS_USERNAME", ""),
|
||||||
RedisPassword: env("REDIS_PASSWORD", ""),
|
RedisPassword: env("REDIS_PASSWORD", ""),
|
||||||
|
|||||||
@@ -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 {
|
type recordingFastAppender struct {
|
||||||
count int
|
count int
|
||||||
batchCount int
|
batchCount int
|
||||||
|
|||||||
Reference in New Issue
Block a user