From 0344c65e5bfec6f3d3ff82797bb4fc60d88b52a8 Mon Sep 17 00:00:00 2001 From: lingniu Date: Fri, 3 Jul 2026 11:59:33 +0800 Subject: [PATCH] feat(go): add nats fast writer --- .../cmd/nats-fast-writer/main.go | 317 ++++++++++++++++++ .../cmd/nats-fast-writer/main_test.go | 73 ++++ go/vehicle-gateway/cmd/realtime-api/main.go | 2 +- go/vehicle-gateway/internal/realtime/model.go | 2 +- .../internal/realtime/repository.go | 37 +- .../internal/realtime/repository_test.go | 41 ++- 6 files changed, 466 insertions(+), 6 deletions(-) create mode 100644 go/vehicle-gateway/cmd/nats-fast-writer/main.go create mode 100644 go/vehicle-gateway/cmd/nats-fast-writer/main_test.go diff --git a/go/vehicle-gateway/cmd/nats-fast-writer/main.go b/go/vehicle-gateway/cmd/nats-fast-writer/main.go new file mode 100644 index 00000000..56f6d41f --- /dev/null +++ b/go/vehicle-gateway/cmd/nats-fast-writer/main.go @@ -0,0 +1,317 @@ +package main + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "log/slog" + "os" + "os/signal" + "strings" + "syscall" + "time" + + "github.com/nats-io/nats.go" + "github.com/redis/go-redis/v9" + _ "github.com/taosdata/driver-go/v3/taosWS" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/health" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/history" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/observability" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/realtime" + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/topics" +) + +func main() { + logger := observability.NewLogger("nats-fast-writer") + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + + cfg := loadConfig() + conn, err := nats.Connect(cfg.NATSURL, nats.Name(cfg.NATSClientName), nats.Timeout(5*time.Second)) + if err != nil { + logger.Error("nats connect failed", "error", err) + os.Exit(1) + } + defer conn.Close() + js, err := conn.JetStream() + if err != nil { + logger.Error("nats jetstream init failed", "error", err) + os.Exit(1) + } + if err := ensureStream(js, cfg); err != nil { + logger.Error("nats stream ensure failed", "stream", cfg.NATSStream, "error", err) + os.Exit(1) + } + sub, err := js.PullSubscribe( + cfg.NATSFilter, + cfg.NATSDurable, + nats.BindStream(cfg.NATSStream), + nats.ManualAck(), + nats.AckWait(cfg.AckWait), + nats.MaxDeliver(-1), + ) + if err != nil { + logger.Error("nats pull consumer init failed", "stream", cfg.NATSStream, "durable", cfg.NATSDurable, "error", err) + os.Exit(1) + } + + tdDB, err := sql.Open(cfg.TDengineDriver, cfg.TDengineDSN) + if err != nil { + logger.Error("tdengine open failed", "error", err) + os.Exit(1) + } + defer tdDB.Close() + tdDB.SetMaxOpenConns(1) + tdDB.SetMaxIdleConns(1) + if err := tdDB.PingContext(ctx); err != nil { + logger.Error("tdengine ping failed", "error", err) + os.Exit(1) + } + historyWriter := history.NewWriter(tdDB) + if cfg.TDengineEnsureSchema { + if err := historyWriter.EnsureSchema(ctx, cfg.TDengineDatabase); err != nil { + logger.Error("tdengine schema bootstrap failed", "error", err) + os.Exit(1) + } + } + if strings.TrimSpace(cfg.TDengineDatabase) != "" { + if _, err := tdDB.ExecContext(ctx, "USE "+cfg.TDengineDatabase); err != nil { + logger.Error("tdengine use database failed", "database", cfg.TDengineDatabase, "error", err) + os.Exit(1) + } + } + + redisClient := redis.NewClient(&redis.Options{ + Addr: cfg.RedisAddr, + Username: cfg.RedisUsername, + Password: cfg.RedisPassword, + DB: cfg.RedisDB, + }) + defer redisClient.Close() + if err := redisClient.Ping(ctx).Err(); err != nil { + logger.Error("redis ping failed", "error", err) + os.Exit(1) + } + realtimeRepo := realtime.NewRepository(redisClient, realtime.Config{OnlineTTL: cfg.OnlineTTL}) + + registry := metrics.NewRegistry() + health.Start(ctx, logger, health.NewServer(env("HEALTH_ADDR", ""), "nats-fast-writer", []health.Check{ + {Name: "nats", Check: func(context.Context) error { + if conn.Status() != nats.CONNECTED { + return fmt.Errorf("nats status is %s", conn.Status().String()) + } + return nil + }}, + {Name: "tdengine", Check: tdDB.PingContext}, + {Name: "redis", Check: func(ctx context.Context) error { return redisClient.Ping(ctx).Err() }}, + }, registry)) + + logger.Info("nats fast writer started", + "stream", cfg.NATSStream, + "durable", cfg.NATSDurable, + "filter", cfg.NATSFilter, + "batch_size", cfg.BatchSize, + "workers", cfg.Workers, + "operation_timeout_ms", cfg.OperationWait.Milliseconds()) + for i := 0; i < cfg.Workers; i++ { + go runFastWorker(ctx, logger, registry, sub, historyWriter, realtimeRepo, cfg) + } + <-ctx.Done() +} + +type config struct { + NATSURL string + NATSClientName string + NATSStream string + NATSDurable string + NATSFilter string + NATSSubjects []string + BatchSize int + FetchWait time.Duration + OperationWait time.Duration + AckWait time.Duration + StreamMaxAge time.Duration + Workers int + TDengineDriver string + TDengineDSN string + TDengineDatabase string + TDengineEnsureSchema bool + RedisAddr string + RedisUsername string + RedisPassword string + RedisDB int + OnlineTTL time.Duration +} + +func loadConfig() config { + subjects := splitCSV(env("NATS_STREAM_SUBJECTS", strings.Join([]string{ + env("NATS_SUBJECT_GB32960_RAW", env("KAFKA_TOPIC_GB32960_RAW", topics.RawGB32960)), + env("NATS_SUBJECT_JT808_RAW", env("KAFKA_TOPIC_JT808_RAW", topics.RawJT808)), + env("NATS_SUBJECT_YUTONG_MQTT_RAW", env("KAFKA_TOPIC_YUTONG_MQTT_RAW", topics.RawYutongMQTT)), + }, ","))) + return config{ + NATSURL: env("NATS_URL", "nats://127.0.0.1:4222"), + NATSClientName: env("NATS_CLIENT_NAME", "lingniu-nats-fast-writer"), + NATSStream: env("NATS_STREAM", "VEHICLE_INGEST"), + NATSDurable: env("NATS_DURABLE", "vehicle-fast-writer"), + NATSFilter: env("NATS_FILTER", "vehicle.raw.go.>"), + 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, + 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), + TDengineDriver: env("TDENGINE_DRIVER", "taosWS"), + TDengineDSN: env("TDENGINE_DSN", ""), + TDengineDatabase: env("TDENGINE_DATABASE", history.DefaultDatabase), + TDengineEnsureSchema: env("TDENGINE_ENSURE_SCHEMA", "true") != "false", + RedisAddr: env("REDIS_ADDR", "127.0.0.1:6379"), + RedisUsername: env("REDIS_USERNAME", ""), + RedisPassword: env("REDIS_PASSWORD", ""), + RedisDB: envInt("REDIS_DB", 0), + OnlineTTL: time.Duration(envInt("REALTIME_ONLINE_TTL_SECONDS", 60)) * time.Second, + } +} + +type fastAppender interface { + AppendAll(context.Context, envelope.FrameEnvelope) error +} + +type fastUpdater interface { + FastUpdate(context.Context, envelope.FrameEnvelope) error +} + +type fastMessage struct { + subject string + data []byte + ack func() error +} + +func runFastWorker(ctx context.Context, logger *slog.Logger, registry *metrics.Registry, sub natsPullSubscription, appender fastAppender, updater fastUpdater, cfg config) { + for { + if ctx.Err() != nil { + return + } + msgs, err := sub.Fetch(cfg.BatchSize, nats.MaxWait(cfg.FetchWait)) + if err != nil { + if errors.Is(err, nats.ErrTimeout) { + continue + } + logger.Error("nats fetch failed", "error", err) + time.Sleep(time.Second) + continue + } + for _, msg := range msgs { + msg := msg + operationCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), cfg.OperationWait) + err := processFastMessage(operationCtx, appender, updater, &fastMessage{ + subject: msg.Subject, + data: msg.Data, + ack: func() error { + return msg.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") + } + } +} + +type natsPullSubscription interface { + Fetch(int, ...nats.PullOpt) ([]*nats.Msg, error) +} + +func processFastMessage(ctx context.Context, appender fastAppender, updater fastUpdater, msg *fastMessage) error { + var env envelope.FrameEnvelope + if err := json.Unmarshal(msg.data, &env); err != nil { + if msg.ack != nil { + _ = msg.ack() + } + return nil + } + if err := appender.AppendAll(ctx, env); err != nil { + return fmt.Errorf("tdengine append: %w", err) + } + if err := updater.FastUpdate(ctx, env); err != nil { + return fmt.Errorf("redis fast update: %w", err) + } + if msg.ack != nil { + if err := msg.ack(); err != nil { + return fmt.Errorf("nats ack: %w", err) + } + } + return nil +} + +func ensureStream(js nats.JetStreamContext, cfg config) error { + stream := &nats.StreamConfig{ + Name: cfg.NATSStream, + Subjects: cfg.NATSSubjects, + Storage: nats.FileStorage, + Retention: nats.LimitsPolicy, + MaxAge: cfg.StreamMaxAge, + Duplicates: 2 * time.Minute, + } + if _, err := js.StreamInfo(cfg.NATSStream); err == nil { + _, err = js.UpdateStream(stream) + return err + } + _, err := js.AddStream(stream) + if err == nil { + return nil + } + _, updateErr := js.UpdateStream(stream) + if updateErr == nil { + return nil + } + return err +} + +func addFastMetric(registry *metrics.Registry, subject string, status string) { + if registry == nil { + return + } + registry.IncCounter("vehicle_fast_writer_messages_total", metrics.Labels{"subject": subject, "status": status}) +} + +func env(key string, fallback string) string { + value := strings.TrimSpace(os.Getenv(key)) + if value == "" { + return fallback + } + return value +} + +func envInt(key string, fallback int) int { + value := strings.TrimSpace(os.Getenv(key)) + if value == "" { + return fallback + } + var parsed int + if _, err := fmt.Sscanf(value, "%d", &parsed); err != nil { + return fallback + } + return parsed +} + +func splitCSV(value string) []string { + var out []string + for _, item := range strings.Split(value, ",") { + item = strings.TrimSpace(item) + if item != "" { + out = append(out, item) + } + } + return out +} diff --git a/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go b/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go new file mode 100644 index 00000000..324cf137 --- /dev/null +++ b/go/vehicle-gateway/cmd/nats-fast-writer/main_test.go @@ -0,0 +1,73 @@ +package main + +import ( + "context" + "errors" + "testing" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +func TestProcessFastMessageWritesTDengineAndRedisBeforeAck(t *testing.T) { + env := envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "VIN001", EventID: "evt-1"} + payload, err := env.MarshalJSONBytes() + if err != nil { + t.Fatal(err) + } + appender := &recordingFastAppender{} + updater := &recordingFastUpdater{} + ackCount := 0 + msg := &fastMessage{data: payload, ack: func() error { + ackCount++ + return nil + }} + + if err := processFastMessage(context.Background(), appender, updater, msg); err != nil { + t.Fatalf("processFastMessage() error = %v", err) + } + if appender.count != 1 || updater.count != 1 || ackCount != 1 { + t.Fatalf("appends=%d updates=%d acks=%d", appender.count, updater.count, ackCount) + } +} + +func TestProcessFastMessageDoesNotAckWhenRedisFails(t *testing.T) { + env := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, VIN: "VIN001", EventID: "evt-2"} + payload, err := env.MarshalJSONBytes() + if err != nil { + t.Fatal(err) + } + appender := &recordingFastAppender{} + updater := &recordingFastUpdater{err: errors.New("redis down")} + ackCount := 0 + msg := &fastMessage{data: payload, ack: func() error { + ackCount++ + return nil + }} + + if err := processFastMessage(context.Background(), appender, updater, msg); err == nil { + t.Fatal("processFastMessage() error = nil, want redis failure") + } + if ackCount != 0 { + t.Fatalf("ack count = %d, want 0", ackCount) + } +} + +type recordingFastAppender struct { + count int + err error +} + +func (a *recordingFastAppender) AppendAll(context.Context, envelope.FrameEnvelope) error { + a.count++ + return a.err +} + +type recordingFastUpdater struct { + count int + err error +} + +func (u *recordingFastUpdater) FastUpdate(context.Context, envelope.FrameEnvelope) error { + u.count++ + return u.err +} diff --git a/go/vehicle-gateway/cmd/realtime-api/main.go b/go/vehicle-gateway/cmd/realtime-api/main.go index 2402a064..7afc75d6 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main.go +++ b/go/vehicle-gateway/cmd/realtime-api/main.go @@ -45,7 +45,7 @@ func main() { } repository := realtime.NewRepository(client, realtime.Config{ - OnlineTTL: time.Duration(envInt("ONLINE_TTL_SECONDS", 600)) * time.Second, + OnlineTTL: time.Duration(envInt("ONLINE_TTL_SECONDS", envInt("REALTIME_ONLINE_TTL_SECONDS", 60))) * time.Second, }) mux := http.NewServeMux() diff --git a/go/vehicle-gateway/internal/realtime/model.go b/go/vehicle-gateway/internal/realtime/model.go index 2414782d..da2a9488 100644 --- a/go/vehicle-gateway/internal/realtime/model.go +++ b/go/vehicle-gateway/internal/realtime/model.go @@ -40,7 +40,7 @@ type Config struct { func (c Config) ttl() time.Duration { if c.OnlineTTL <= 0 { - return 10 * time.Minute + return time.Minute } return c.OnlineTTL } diff --git a/go/vehicle-gateway/internal/realtime/repository.go b/go/vehicle-gateway/internal/realtime/repository.go index b93e8511..54e00c7d 100644 --- a/go/vehicle-gateway/internal/realtime/repository.go +++ b/go/vehicle-gateway/internal/realtime/repository.go @@ -27,6 +27,21 @@ func NewRepository(client *redis.Client, cfg Config) *Repository { return &Repository{client: client, cfg: cfg} } +func (r *Repository) FastUpdate(ctx context.Context, env envelope.FrameEnvelope) error { + vin := strings.TrimSpace(env.VIN) + if vin == "" { + return nil + } + vehicleKey := strings.TrimSpace(env.VehicleKey()) + if vehicleKey == "" || strings.HasSuffix(vehicleKey, ":unknown") { + return nil + } + if err := r.setKV(ctx, vin, env, env.Parsed); err != nil { + return err + } + return r.setOnline(ctx, vehicleKey, vin, []envelope.Protocol{env.Protocol}, env.ReceivedAtMS) +} + func (r *Repository) Update(ctx context.Context, env envelope.FrameEnvelope) error { vin := strings.TrimSpace(env.VIN) if vin == "" { @@ -124,7 +139,7 @@ func (r *Repository) Update(ctx context.Context, env envelope.FrameEnvelope) err Protocols: protocols, TTLSeconds: int64(r.cfg.ttl().Seconds()), } - if err := r.setJSON(ctx, onlineKey(vehicleKey), online, r.cfg.ttl()); err != nil { + if err := r.setOnlineStatus(ctx, online); err != nil { return err } return r.client.ZAdd(ctx, "vehicle:last_seen", redis.Z{Score: float64(env.ReceivedAtMS), Member: vehicleKey}).Err() @@ -209,12 +224,30 @@ func (r *Repository) setKV(ctx context.Context, vin string, env envelope.FrameEn pipe := r.client.Pipeline() for key, values := range grouped { pipe.HSet(ctx, key, values) - pipe.Expire(ctx, key, r.cfg.ttl()) } _, err := pipe.Exec(ctx) return err } +func (r *Repository) setOnline(ctx context.Context, vehicleKey string, vin string, protocols []envelope.Protocol, lastSeenMS int64) error { + online := OnlineStatus{ + VehicleKey: vehicleKey, + VIN: vin, + Online: true, + LastSeenMS: lastSeenMS, + Protocols: protocols, + TTLSeconds: int64(r.cfg.ttl().Seconds()), + } + if err := r.setOnlineStatus(ctx, online); err != nil { + return err + } + return r.client.ZAdd(ctx, "vehicle:last_seen", redis.Z{Score: float64(lastSeenMS), Member: vehicleKey}).Err() +} + +func (r *Repository) setOnlineStatus(ctx context.Context, online OnlineStatus) error { + return r.setJSON(ctx, onlineKey(online.VehicleKey), online, r.cfg.ttl()) +} + func (r *Repository) addProtocol(ctx context.Context, vehicleKey string, protocol envelope.Protocol) ([]envelope.Protocol, error) { key := protocolsKey(vehicleKey) if err := r.client.SAdd(ctx, key, string(protocol)).Err(); err != nil { diff --git a/go/vehicle-gateway/internal/realtime/repository_test.go b/go/vehicle-gateway/internal/realtime/repository_test.go index f3e65cf5..d359f1d5 100644 --- a/go/vehicle-gateway/internal/realtime/repository_test.go +++ b/go/vehicle-gateway/internal/realtime/repository_test.go @@ -308,8 +308,45 @@ func TestRepositoryWritesRealtimeKVHashes(t *testing.T) { if stackKV["stack_water_outlet_temp_c"] != "63" { t.Fatalf("stack kv = %#v", stackKV) } - if ttl := repo.client.TTL(ctx, "vehicle:rt-kv:GB32960:VIN001:vehicle").Val(); ttl <= 0 { - t.Fatalf("kv ttl should be set, got %v", ttl) + if ttl := repo.client.TTL(ctx, "vehicle:rt-kv:GB32960:VIN001:vehicle").Val(); ttl != -1 { + t.Fatalf("kv ttl should not expire, got %v", ttl) + } +} + +func TestRepositoryFastUpdateOnlyWritesPermanentKVAndMinuteOnline(t *testing.T) { + repo, closeFn := newTestRepository(t) + defer closeFn() + ctx := context.Background() + + if err := repo.FastUpdate(ctx, envelope.FrameEnvelope{ + Protocol: envelope.ProtocolGB32960, + VIN: "VIN001", + EventTimeMS: 1000, + ReceivedAtMS: 1100, + Parsed: map[string]any{ + "data_units": []any{ + map[string]any{"type": "0x01", "name": "vehicle", "value": map[string]any{"soc_percent": 88.0}}, + }, + }, + }); err != nil { + t.Fatalf("FastUpdate() error = %v", err) + } + + vehicleKV, err := repo.client.HGetAll(ctx, "vehicle:rt-kv:GB32960:VIN001:vehicle").Result() + if err != nil { + t.Fatalf("vehicle kv HGetAll error = %v", err) + } + if vehicleKV["soc_percent"] != "88" || vehicleKV["_event_time_ms"] != "1000" { + t.Fatalf("vehicle kv = %#v", vehicleKV) + } + if ttl := repo.client.TTL(ctx, "vehicle:rt-kv:GB32960:VIN001:vehicle").Val(); ttl != -1 { + t.Fatalf("kv ttl should not expire, got %v", ttl) + } + if ttl := repo.client.TTL(ctx, "vehicle:online:VIN001").Val(); ttl <= 0 || ttl > time.Minute { + t.Fatalf("online ttl = %v, want within 1 minute", ttl) + } + if repo.client.Exists(ctx, "vehicle:latest:VIN001").Val() != 0 { + t.Fatal("fast update should not write merged snapshot") } }