diff --git a/docs/architecture/production-data-plane-inventory.md b/docs/architecture/production-data-plane-inventory.md index d9b7df5a..7f6b40bd 100644 --- a/docs/architecture/production-data-plane-inventory.md +++ b/docs/architecture/production-data-plane-inventory.md @@ -28,6 +28,7 @@ NATS 部署在 Kafka ECS 内网 `172.17.111.56:4222`。 | Subjects | `vehicle.raw.go.gb32960.v1`、`vehicle.raw.go.jt808.v1`、`vehicle.raw.go.yutong-mqtt.v1` | | Durable consumer | `vehicle-kafka-bridge` | | 语义 | Gateway 先写 NATS raw subject;bridge 写 Kafka 成功后才 ACK NATS | +| 保留策略 | `NATS_STREAM_MAX_BYTES=21474836480`,`NATS_STREAM_MAX_AGE_HOURS=24`,`NATS_STREAM_ENSURE_TIMEOUT_SECONDS=60`,NATS 仅做入口缓冲 | ### Kafka diff --git a/docs/ops/100k-capacity-baseline.md b/docs/ops/100k-capacity-baseline.md index 13a529a4..ed7af7bb 100644 --- a/docs/ops/100k-capacity-baseline.md +++ b/docs/ops/100k-capacity-baseline.md @@ -88,6 +88,19 @@ systemd gateway 已配置 `LimitNOFILE=1048576`,需要持续保持。 | Realtime MySQL | async queue depth | 不持续增长 | | Stat writer | `vehicle_stat_write_duration_ms_histogram_bucket` | p99 不持续上升 | +## Retention Guardrails + +10W 车辆接入时,中间件只承担解耦和短期缓冲职责,不能成为长期历史库。 + +| Component | Guardrail | +| --- | --- | +| NATS JetStream `VEHICLE_INGEST` | `NATS_STREAM_MAX_BYTES=21474836480`,`NATS_STREAM_MAX_AGE_HOURS=24`,`NATS_STREAM_ENSURE_TIMEOUT_SECONDS=60` | +| Kafka `vehicle.raw.go.*` / `vehicle.fields.go.*` | `retention.ms=21600000`、`segment.ms=600000`、`segment.bytes=268435456` | +| Fast writer | `FAST_WRITER_OPERATION_TIMEOUT_MS=1000`,避免 TDengine 批量写尾延迟导致整批重投 | +| Root disk | 使用率低于 `80%`,超过 `85%` 进入容量告警 | + +如果 NATS `consumer_pending` 下降但仍高,说明系统在追历史积压;如果 `ack_pending=0` 且 Kafka lag 为 `0`,不要重启服务,重点观察 NATS data size 和 pending 下降斜率。 + ## 分阶段压测 压测入口: diff --git a/docs/ops/vehicle-ingest-runbook.md b/docs/ops/vehicle-ingest-runbook.md index e1b41bd9..b0f5726d 100644 --- a/docs/ops/vehicle-ingest-runbook.md +++ b/docs/ops/vehicle-ingest-runbook.md @@ -196,13 +196,30 @@ Kafka 是可回放日志。只要 Kafka 还保留消息,恢复 consumer 后 la 1. 检查 bridge `/readyz` 和日志。 2. 对比 `vehicle_bridge_kafka_writes_total` 和 `vehicle_bridge_nats_acks_total`。 3. 如果 Kafka 写入失败,检查 ECS 到 Kafka broker 的网络。 -4. 如果 ack-pending 卡住且日志持续报错,只重启 bridge。 +4. 如果 `ack_pending > 0`,优先检查 Kafka 是否监听 `9092`、Kafka 盘是否打满,再看 bridge 日志。 +5. 如果 `ack_pending = 0` 但 `consumer_pending` 下降,说明 bridge 正在追历史积压;不要反复重启,持续观察下降速度即可。 +6. 如果 `consumer_pending` 不下降,才考虑扩容 bridge 或降低 batch/fetch 等配置。 ```bash journalctl -u lingniu-go-nats-kafka-bridge.service --since '10 minutes ago' --no-pager systemctl restart lingniu-go-nats-kafka-bridge.service ``` +生产 stream 必须设置字节上限,避免 NATS JetStream 在 Kafka 或下游故障时吃满根盘: + +```bash +grep NATS_STREAM_MAX_BYTES /opt/lingniu-go-native/env/nats-fast-writer.env +grep NATS_STREAM_MAX_BYTES /opt/lingniu-go-native/env/nats-kafka-bridge.env +grep NATS_STREAM_ENSURE_TIMEOUT_SECONDS /opt/lingniu-go-native/env/nats-fast-writer.env +grep NATS_STREAM_ENSURE_TIMEOUT_SECONDS /opt/lingniu-go-native/env/nats-kafka-bridge.env +du -sh /opt/lingniu-nats/data +df -h / +``` + +当前建议值:`NATS_STREAM_MAX_BYTES=21474836480`,即 `20GiB`;`NATS_STREAM_ENSURE_TIMEOUT_SECONDS=60`,避免大 stream 元数据更新时被 NATS 客户端默认 5s 超时误杀。Kafka topic 只作为短期缓冲,当前建议保留 `6h`,不要把 Kafka 或 NATS 当长期历史存储;长期历史和 RAW 查询以 TDengine/MySQL 投影为准。 + +fast-writer 的 `FAST_WRITER_OPERATION_TIMEOUT_MS` 建议为 `1000`。实时链路仍以 100ms 级为目标,但 TDengine 批量写存在尾延迟,过小的超时会造成 NATS 消息反复重投和重复写压力。 + ### Raw 有数据但实时查不到 1. 查 realtime 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 969fb3ba..eb5afb39 100644 --- a/go/vehicle-gateway/cmd/nats-fast-writer/main.go +++ b/go/vehicle-gateway/cmd/nats-fast-writer/main.go @@ -132,6 +132,8 @@ type config struct { OperationWait time.Duration AckWait time.Duration StreamMaxAge time.Duration + StreamMaxBytes int64 + StreamEnsureWait time.Duration Workers int TDengineDriver string TDengineDSN string @@ -167,9 +169,11 @@ func loadConfig() config { NATSSubjects: subjects, BatchSize: envInt("FAST_WRITER_BATCH_SIZE", 100), FetchWait: time.Duration(envInt("FAST_WRITER_FETCH_WAIT_MS", 100)) * time.Millisecond, - OperationWait: time.Duration(envInt("FAST_WRITER_OPERATION_TIMEOUT_MS", 100)) * time.Millisecond, + OperationWait: time.Duration(envInt("FAST_WRITER_OPERATION_TIMEOUT_MS", 1000)) * 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, + StreamMaxBytes: envInt64("NATS_STREAM_MAX_BYTES", 20*1024*1024*1024), + StreamEnsureWait: time.Duration(envInt("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", 60)) * time.Second, Workers: workers, TDengineDriver: env("TDENGINE_DRIVER", "taosWS"), TDengineDSN: env("TDENGINE_DSN", ""), @@ -381,23 +385,27 @@ func fastBatchSubject(messages []*fastMessage) string { } func ensureStream(js nats.JetStreamContext, cfg config) error { + ctx, cancel := context.WithTimeout(context.Background(), cfg.StreamEnsureWait) + defer cancel() + opts := []nats.JSOpt{nats.Context(ctx)} stream := &nats.StreamConfig{ Name: cfg.NATSStream, Subjects: cfg.NATSSubjects, Storage: nats.FileStorage, Retention: nats.LimitsPolicy, MaxAge: cfg.StreamMaxAge, + MaxBytes: cfg.StreamMaxBytes, Duplicates: 2 * time.Minute, } - if _, err := js.StreamInfo(cfg.NATSStream); err == nil { - _, err = js.UpdateStream(stream) + if _, err := js.StreamInfo(cfg.NATSStream, opts...); err == nil { + _, err = js.UpdateStream(stream, opts...) return err } - _, err := js.AddStream(stream) + _, err := js.AddStream(stream, opts...) if err == nil { return nil } - _, updateErr := js.UpdateStream(stream) + _, updateErr := js.UpdateStream(stream, opts...) if updateErr == nil { return nil } @@ -467,6 +475,18 @@ func envInt(key string, fallback int) int { return parsed } +func envInt64(key string, fallback int64) int64 { + value := strings.TrimSpace(os.Getenv(key)) + if value == "" { + return fallback + } + var parsed int64 + if _, err := fmt.Sscanf(value, "%d", &parsed); err != nil || parsed <= 0 { + return fallback + } + return parsed +} + func splitCSV(value string) []string { var out []string for _, item := range strings.Split(value, ",") { 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 b4e67e67..cdf15c65 100644 --- a/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go +++ b/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go @@ -6,6 +6,7 @@ import ( "strings" "sync" "testing" + "time" "github.com/nats-io/nats.go" @@ -243,6 +244,38 @@ func TestLoadConfigDefaultsTDenginePoolToSingleConnection(t *testing.T) { } } +func TestLoadConfigDefaultsStreamMaxBytes(t *testing.T) { + cfg := loadConfig() + + if got, want := cfg.StreamMaxBytes, int64(20*1024*1024*1024); got != want { + t.Fatalf("StreamMaxBytes = %d, want %d", got, want) + } + if got, want := cfg.StreamEnsureWait, 60*time.Second; got != want { + t.Fatalf("StreamEnsureWait = %v, want %v", got, want) + } + if got, want := cfg.OperationWait, time.Second; got != want { + t.Fatalf("OperationWait = %v, want %v", got, want) + } +} + +func TestLoadConfigReadsStreamMaxBytesOverride(t *testing.T) { + t.Setenv("NATS_STREAM_MAX_BYTES", "1073741824") + t.Setenv("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", "90") + t.Setenv("FAST_WRITER_OPERATION_TIMEOUT_MS", "1500") + + cfg := loadConfig() + + if got, want := cfg.StreamMaxBytes, int64(1073741824); got != want { + t.Fatalf("StreamMaxBytes = %d, want %d", got, want) + } + if got, want := cfg.StreamEnsureWait, 90*time.Second; got != want { + t.Fatalf("StreamEnsureWait = %v, want %v", got, want) + } + if got, want := cfg.OperationWait, 1500*time.Millisecond; got != want { + t.Fatalf("OperationWait = %v, want %v", got, want) + } +} + func TestLoadConfigReadsTDenginePoolOverride(t *testing.T) { t.Setenv("FAST_WRITER_WORKERS", "12") t.Setenv("FAST_WRITER_TDENGINE_MAX_OPEN_CONNS", "4") diff --git a/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go b/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go index 901581d6..1dfdb67c 100644 --- a/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go +++ b/go/vehicle-gateway/cmd/nats-kafka-bridge/main.go @@ -85,19 +85,21 @@ func main() { } type config struct { - NATSURL string - NATSClientName string - NATSStream string - NATSDurable string - NATSFilter string - NATSSubjects []string - KafkaBrokers []string - Route map[string]string - BatchSize int - FetchWait time.Duration - OperationWait time.Duration - AckWait time.Duration - StreamMaxAge time.Duration + NATSURL string + NATSClientName string + NATSStream string + NATSDurable string + NATSFilter string + NATSSubjects []string + KafkaBrokers []string + Route map[string]string + BatchSize int + FetchWait time.Duration + OperationWait time.Duration + AckWait time.Duration + StreamMaxAge time.Duration + StreamMaxBytes int64 + StreamEnsureWait time.Duration } func loadConfig() config { @@ -113,19 +115,21 @@ func loadConfig() config { route[unifiedSubject] = unifiedTopic } return config{ - NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"), - NATSClientName: env("NATS_CLIENT_NAME", "lingniu-nats-kafka-bridge"), - NATSStream: env("NATS_STREAM", "VEHICLE_INGEST"), - NATSDurable: env("NATS_DURABLE", "vehicle-kafka-bridge"), - NATSFilter: env("NATS_FILTER", "vehicle.>"), - NATSSubjects: splitCSV(env("NATS_STREAM_SUBJECTS", strings.Join(mapKeys(route), ","))), - KafkaBrokers: splitCSV(env("KAFKA_BROKERS", "127.0.0.1:9092")), - Route: route, - BatchSize: envInt("BRIDGE_BATCH_SIZE", 500), - FetchWait: time.Duration(envInt("BRIDGE_FETCH_WAIT_MS", 1000)) * time.Millisecond, - OperationWait: time.Duration(envInt("BRIDGE_OPERATION_TIMEOUT_MS", 30000)) * time.Millisecond, - AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 60)) * time.Second, - StreamMaxAge: time.Duration(envInt("NATS_STREAM_MAX_AGE_HOURS", 24)) * time.Hour, + NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"), + NATSClientName: env("NATS_CLIENT_NAME", "lingniu-nats-kafka-bridge"), + NATSStream: env("NATS_STREAM", "VEHICLE_INGEST"), + NATSDurable: env("NATS_DURABLE", "vehicle-kafka-bridge"), + NATSFilter: env("NATS_FILTER", "vehicle.>"), + NATSSubjects: splitCSV(env("NATS_STREAM_SUBJECTS", strings.Join(mapKeys(route), ","))), + KafkaBrokers: splitCSV(env("KAFKA_BROKERS", "127.0.0.1:9092")), + Route: route, + BatchSize: envInt("BRIDGE_BATCH_SIZE", 500), + FetchWait: time.Duration(envInt("BRIDGE_FETCH_WAIT_MS", 1000)) * time.Millisecond, + OperationWait: time.Duration(envInt("BRIDGE_OPERATION_TIMEOUT_MS", 30000)) * time.Millisecond, + AckWait: time.Duration(envInt("NATS_ACK_WAIT_SECONDS", 60)) * time.Second, + StreamMaxAge: time.Duration(envInt("NATS_STREAM_MAX_AGE_HOURS", 24)) * time.Hour, + StreamMaxBytes: envInt64("NATS_STREAM_MAX_BYTES", 20*1024*1024*1024), + StreamEnsureWait: time.Duration(envInt("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", 60)) * time.Second, } } @@ -213,23 +217,27 @@ func runBridge(ctx context.Context, logger *slog.Logger, registry *metrics.Regis } func ensureStream(js nats.JetStreamContext, cfg config) error { + ctx, cancel := context.WithTimeout(context.Background(), cfg.StreamEnsureWait) + defer cancel() + opts := []nats.JSOpt{nats.Context(ctx)} stream := &nats.StreamConfig{ Name: cfg.NATSStream, Subjects: cfg.NATSSubjects, Storage: nats.FileStorage, Retention: nats.LimitsPolicy, MaxAge: cfg.StreamMaxAge, + MaxBytes: cfg.StreamMaxBytes, Duplicates: 2 * time.Minute, } - if _, err := js.StreamInfo(cfg.NATSStream); err == nil { - _, err = js.UpdateStream(stream) + if _, err := js.StreamInfo(cfg.NATSStream, opts...); err == nil { + _, err = js.UpdateStream(stream, opts...) return err } - _, err := js.AddStream(stream) + _, err := js.AddStream(stream, opts...) if err == nil { return nil } - _, updateErr := js.UpdateStream(stream) + _, updateErr := js.UpdateStream(stream, opts...) if updateErr == nil { return nil } @@ -362,6 +370,18 @@ func envInt(key string, fallback int) int { return parsed } +func envInt64(key string, fallback int64) int64 { + value := strings.TrimSpace(os.Getenv(key)) + if value == "" { + return fallback + } + parsed, err := strconv.ParseInt(value, 10, 64) + if err != nil || parsed <= 0 { + return fallback + } + return parsed +} + func splitCSV(value string) []string { var out []string for _, item := range strings.Split(value, ",") { 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 3356f966..e8ebe393 100644 --- a/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go +++ b/go/vehicle-gateway/cmd/nats-kafka-bridge/main_test.go @@ -6,6 +6,7 @@ import ( "errors" "strings" "testing" + "time" "github.com/nats-io/nats.go" "github.com/segmentio/kafka-go" @@ -177,6 +178,31 @@ func TestLoadConfigDefaultsToGoSubjectRoutes(t *testing.T) { } } +func TestLoadConfigDefaultsStreamMaxBytes(t *testing.T) { + cfg := loadConfig() + + if got, want := cfg.StreamMaxBytes, int64(20*1024*1024*1024); got != want { + t.Fatalf("StreamMaxBytes = %d, want %d", got, want) + } + if got, want := cfg.StreamEnsureWait, 60*time.Second; got != want { + t.Fatalf("StreamEnsureWait = %v, want %v", got, want) + } +} + +func TestLoadConfigReadsStreamMaxBytesOverride(t *testing.T) { + t.Setenv("NATS_STREAM_MAX_BYTES", "1073741824") + t.Setenv("NATS_STREAM_ENSURE_TIMEOUT_SECONDS", "90") + + cfg := loadConfig() + + if got, want := cfg.StreamMaxBytes, int64(1073741824); got != want { + t.Fatalf("StreamMaxBytes = %d, want %d", got, want) + } + if got, want := cfg.StreamEnsureWait, 90*time.Second; got != want { + t.Fatalf("StreamEnsureWait = %v, want %v", got, want) + } +} + func TestLoadConfigIncludesUnifiedOnlyWhenExplicitlyConfigured(t *testing.T) { t.Setenv("NATS_SUBJECT_UNIFIED", "vehicle.event.go.unified.v1") t.Setenv("KAFKA_TOPIC_UNIFIED", "vehicle.event.go.unified.v1")