perf(go): cache realtime plate lookups

This commit is contained in:
lingniu
2026-07-03 07:42:24 +08:00
parent 9345fe60b5
commit 1c0d7dc03e
5 changed files with 124 additions and 2 deletions

View File

@@ -61,6 +61,7 @@
6. MySQL realtime 当前态表不保留代理主键和创建时间。 6. MySQL realtime 当前态表不保留代理主键和创建时间。
- 原因:`vehicle_realtime_snapshot``vehicle_realtime_location` 都是每协议每 VIN 一行的当前态投影,业务主键就是 `(protocol, vin)` - 原因:`vehicle_realtime_snapshot``vehicle_realtime_location` 都是每协议每 VIN 一行的当前态投影,业务主键就是 `(protocol, vin)`
- 状态schema 使用 `PRIMARY KEY (protocol, vin)`,不保留自增 `id``created_at`,对外只暴露最新 `updated_at` - 状态schema 使用 `PRIMARY KEY (protocol, vin)`,不保留自增 `id``created_at`,对外只暴露最新 `updated_at`
- 写入:车牌从 `vehicle_identity_binding` 反查时使用 VIN 级短 TTL 内存缓存,避免实时帧每条都打 MySQL缓存不作为事实来源。
- 生产:项目上线前可直接重建这两张 MySQL 表,实时数据会从 Kafka 新消息继续投影。 - 生产:项目上线前可直接重建这两张 MySQL 表,实时数据会从 Kafka 新消息继续投影。
7. identity 层只保留两张表。 7. identity 层只保留两张表。

View File

@@ -73,7 +73,11 @@ func main() {
healthChecks = append(healthChecks, health.Check{Name: "mysql", Check: db.PingContext}) healthChecks = append(healthChecks, health.Check{Name: "mysql", Check: db.PingContext})
closeStats = func() { _ = db.Close() } closeStats = func() { _ = db.Close() }
bindingTable := env("VEHICLE_IDENTITY_TABLE", "vehicle_identity_binding") 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 env("MYSQL_REALTIME_SNAPSHOT_ENABLED", "true") != "false" {
if err := snapshotWriter.EnsureSchema(ctx); err != nil { if err := snapshotWriter.EnsureSchema(ctx); err != nil {
_ = db.Close() _ = db.Close()
@@ -81,7 +85,7 @@ func main() {
os.Exit(1) os.Exit(1)
} }
updater = compositeRealtimeUpdater{primary: repository, secondary: snapshotWriter} 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/stats/daily-metrics", stats.NewMetricHandler(stats.NewMetricRepository(db)))
mux.Handle("/api/realtime/snapshots", realtime.NewSnapshotQueryHandler(realtime.NewSnapshotQueryRepository(db))) mux.Handle("/api/realtime/snapshots", realtime.NewSnapshotQueryHandler(realtime.NewSnapshotQueryRepository(db)))

View File

@@ -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 { type contextCheckingRealtimeUpdater struct {
ctxErr error ctxErr error
count int count int

View File

@@ -5,6 +5,7 @@ import (
"database/sql" "database/sql"
"errors" "errors"
"strings" "strings"
"sync"
"time" "time"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
@@ -18,6 +19,69 @@ type PlateResolver interface {
PlateByVIN(context.Context, string) (string, error) 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 { type SnapshotWriter struct {
exec SnapshotExecer exec SnapshotExecer
plateResolver PlateResolver plateResolver PlateResolver

View File

@@ -6,6 +6,7 @@ import (
"errors" "errors"
"strings" "strings"
"testing" "testing"
"time"
"github.com/DATA-DOG/go-sqlmock" "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) { func TestSnapshotWriterKeepsEventPlateWhenPresent(t *testing.T) {
exec := &recordingSnapshotExec{} exec := &recordingSnapshotExec{}
resolver := &recordingPlateResolver{plate: "沪B99999"} resolver := &recordingPlateResolver{plate: "沪B99999"}
@@ -358,9 +397,11 @@ type recordingPlateResolver struct {
vin string vin string
plate string plate string
err error err error
calls int
} }
func (r *recordingPlateResolver) PlateByVIN(_ context.Context, vin string) (string, error) { func (r *recordingPlateResolver) PlateByVIN(_ context.Context, vin string) (string, error) {
r.calls++
r.vin = vin r.vin = vin
return r.plate, r.err return r.plate, r.err
} }