From 1c0d7dc03edda501d789b77abd460c0355ac0443 Mon Sep 17 00:00:00 2001 From: lingniu Date: Fri, 3 Jul 2026 07:42:24 +0800 Subject: [PATCH] perf(go): cache realtime plate lookups --- docs/architecture/storage-minimal-contract.md | 1 + go/vehicle-gateway/cmd/realtime-api/main.go | 8 ++- .../cmd/realtime-api/main_test.go | 12 ++++ .../internal/realtime/snapshot_writer.go | 64 +++++++++++++++++++ .../internal/realtime/snapshot_writer_test.go | 41 ++++++++++++ 5 files changed, 124 insertions(+), 2 deletions(-) diff --git a/docs/architecture/storage-minimal-contract.md b/docs/architecture/storage-minimal-contract.md index 95d5664e..86754f85 100644 --- a/docs/architecture/storage-minimal-contract.md +++ b/docs/architecture/storage-minimal-contract.md @@ -61,6 +61,7 @@ 6. MySQL realtime 当前态表不保留代理主键和创建时间。 - 原因:`vehicle_realtime_snapshot`、`vehicle_realtime_location` 都是每协议每 VIN 一行的当前态投影,业务主键就是 `(protocol, vin)`。 - 状态:schema 使用 `PRIMARY KEY (protocol, vin)`,不保留自增 `id` 和 `created_at`,对外只暴露最新 `updated_at`。 + - 写入:车牌从 `vehicle_identity_binding` 反查时使用 VIN 级短 TTL 内存缓存,避免实时帧每条都打 MySQL;缓存不作为事实来源。 - 生产:项目上线前可直接重建这两张 MySQL 表,实时数据会从 Kafka 新消息继续投影。 7. identity 层只保留两张表。 diff --git a/go/vehicle-gateway/cmd/realtime-api/main.go b/go/vehicle-gateway/cmd/realtime-api/main.go index ee27a276..848258dc 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main.go +++ b/go/vehicle-gateway/cmd/realtime-api/main.go @@ -73,7 +73,11 @@ func main() { healthChecks = append(healthChecks, health.Check{Name: "mysql", Check: db.PingContext}) closeStats = func() { _ = db.Close() } bindingTable := env("VEHICLE_IDENTITY_TABLE", "vehicle_identity_binding") - snapshotWriter := realtime.NewSnapshotWriterWithPlateResolver(db, realtime.NewBindingPlateResolver(db, bindingTable)) + plateResolver := realtime.NewCachedPlateResolver( + realtime.NewBindingPlateResolver(db, bindingTable), + time.Duration(envInt("PLATE_CACHE_TTL_SECONDS", 600))*time.Second, + ) + snapshotWriter := realtime.NewSnapshotWriterWithPlateResolver(db, plateResolver) if env("MYSQL_REALTIME_SNAPSHOT_ENABLED", "true") != "false" { if err := snapshotWriter.EnsureSchema(ctx); err != nil { _ = db.Close() @@ -81,7 +85,7 @@ func main() { os.Exit(1) } updater = compositeRealtimeUpdater{primary: repository, secondary: snapshotWriter} - logger.Info("realtime mysql snapshot enabled", "table", "vehicle_realtime_snapshot", "plate_binding_table", bindingTable) + logger.Info("realtime mysql snapshot enabled", "table", "vehicle_realtime_snapshot", "plate_binding_table", bindingTable, "plate_cache_ttl_seconds", envInt("PLATE_CACHE_TTL_SECONDS", 600)) } mux.Handle("/api/stats/daily-metrics", stats.NewMetricHandler(stats.NewMetricRepository(db))) mux.Handle("/api/realtime/snapshots", realtime.NewSnapshotQueryHandler(realtime.NewSnapshotQueryRepository(db))) diff --git a/go/vehicle-gateway/cmd/realtime-api/main_test.go b/go/vehicle-gateway/cmd/realtime-api/main_test.go index 81970479..3d240a81 100644 --- a/go/vehicle-gateway/cmd/realtime-api/main_test.go +++ b/go/vehicle-gateway/cmd/realtime-api/main_test.go @@ -114,6 +114,18 @@ func TestRealtimeAPIDoesNotExposeDuplicateMileagePointRoute(t *testing.T) { } } +func TestRealtimeAPICachesBindingPlateResolver(t *testing.T) { + source, err := os.ReadFile("main.go") + if err != nil { + t.Fatalf("read main.go: %v", err) + } + for _, want := range []string{"NewCachedPlateResolver", "PLATE_CACHE_TTL_SECONDS"} { + if !strings.Contains(string(source), want) { + t.Fatalf("realtime api should configure cached plate resolver, missing %s", want) + } + } +} + type contextCheckingRealtimeUpdater struct { ctxErr error count int diff --git a/go/vehicle-gateway/internal/realtime/snapshot_writer.go b/go/vehicle-gateway/internal/realtime/snapshot_writer.go index ada9dad6..3b30a458 100644 --- a/go/vehicle-gateway/internal/realtime/snapshot_writer.go +++ b/go/vehicle-gateway/internal/realtime/snapshot_writer.go @@ -5,6 +5,7 @@ import ( "database/sql" "errors" "strings" + "sync" "time" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" @@ -18,6 +19,69 @@ type PlateResolver interface { PlateByVIN(context.Context, string) (string, error) } +type CachedPlateResolver struct { + delegate PlateResolver + ttl time.Duration + now func() time.Time + mu sync.Mutex + entries map[string]cachedPlateEntry +} + +type cachedPlateEntry struct { + plate string + notFound bool + expiresAt time.Time +} + +func NewCachedPlateResolver(delegate PlateResolver, ttl time.Duration) *CachedPlateResolver { + if delegate == nil { + panic("cached plate resolver delegate must not be nil") + } + if ttl <= 0 { + ttl = 10 * time.Minute + } + return &CachedPlateResolver{ + delegate: delegate, + ttl: ttl, + now: time.Now, + entries: map[string]cachedPlateEntry{}, + } +} + +func (r *CachedPlateResolver) PlateByVIN(ctx context.Context, vin string) (string, error) { + vin = strings.TrimSpace(vin) + if vin == "" { + return "", sql.ErrNoRows + } + now := r.now() + r.mu.Lock() + entry, ok := r.entries[vin] + if ok && now.Before(entry.expiresAt) { + r.mu.Unlock() + if entry.notFound { + return "", sql.ErrNoRows + } + return entry.plate, nil + } + r.mu.Unlock() + + plate, err := r.delegate.PlateByVIN(ctx, vin) + if err != nil && !errors.Is(err, sql.ErrNoRows) { + return "", err + } + r.mu.Lock() + r.entries[vin] = cachedPlateEntry{ + plate: strings.TrimSpace(plate), + notFound: errors.Is(err, sql.ErrNoRows), + expiresAt: now.Add(r.ttl), + } + r.mu.Unlock() + if err != nil { + return "", err + } + return strings.TrimSpace(plate), nil +} + type SnapshotWriter struct { exec SnapshotExecer plateResolver PlateResolver diff --git a/go/vehicle-gateway/internal/realtime/snapshot_writer_test.go b/go/vehicle-gateway/internal/realtime/snapshot_writer_test.go index 74e92555..6110e411 100644 --- a/go/vehicle-gateway/internal/realtime/snapshot_writer_test.go +++ b/go/vehicle-gateway/internal/realtime/snapshot_writer_test.go @@ -6,6 +6,7 @@ import ( "errors" "strings" "testing" + "time" "github.com/DATA-DOG/go-sqlmock" @@ -212,6 +213,44 @@ func TestSnapshotWriterBackfillsPlateFromBindingByVIN(t *testing.T) { } } +func TestSnapshotWriterCachesBindingPlateByVIN(t *testing.T) { + exec := &recordingSnapshotExec{} + resolver := &recordingPlateResolver{plate: "沪A12345"} + writer := NewSnapshotWriterWithPlateResolver(exec, NewCachedPlateResolver(resolver, time.Hour)) + event := envelope.FrameEnvelope{ + Protocol: envelope.ProtocolGB32960, + MessageID: "0x02", + VIN: "VIN001", + EventTimeMS: 1782918600000, + ReceivedAtMS: 1782918601000, + Fields: map[string]any{ + envelope.FieldLatitude: 30.590151, + envelope.FieldLongitude: 121.069881, + }, + } + + if err := writer.Update(context.Background(), event); err != nil { + t.Fatalf("first Update() error = %v", err) + } + event.EventTimeMS += 1000 + event.ReceivedAtMS += 1000 + if err := writer.Update(context.Background(), event); err != nil { + t.Fatalf("second Update() error = %v", err) + } + + if resolver.calls != 1 { + t.Fatalf("plate resolver calls = %d, want 1", resolver.calls) + } + if len(exec.calls) != 4 { + t.Fatalf("exec calls = %d, want 4", len(exec.calls)) + } + for _, index := range []int{0, 2} { + if got, want := exec.calls[index].args[2], "沪A12345"; got != want { + t.Fatalf("snapshot %d plate arg = %#v, want %q", index, got, want) + } + } +} + func TestSnapshotWriterKeepsEventPlateWhenPresent(t *testing.T) { exec := &recordingSnapshotExec{} resolver := &recordingPlateResolver{plate: "沪B99999"} @@ -358,9 +397,11 @@ type recordingPlateResolver struct { vin string plate string err error + calls int } func (r *recordingPlateResolver) PlateByVIN(_ context.Context, vin string) (string, error) { + r.calls++ r.vin = vin return r.plate, r.err }