feat(go): add nats fast writer

This commit is contained in:
lingniu
2026-07-03 11:59:33 +08:00
parent 76f5a20b79
commit 0344c65e5b
6 changed files with 466 additions and 6 deletions

View File

@@ -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
}

View File

@@ -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
}

View File

@@ -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()

View File

@@ -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
}

View File

@@ -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 {

View File

@@ -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")
}
}