diff --git a/docs/architecture/iot-data-platform-principles.md b/docs/architecture/iot-data-platform-principles.md index 583c51d2..d9a8dd30 100644 --- a/docs/architecture/iot-data-platform-principles.md +++ b/docs/architecture/iot-data-platform-principles.md @@ -52,7 +52,7 @@ flowchart LR | MySQL | `vehicle_daily_mileage` | Queryable daily mileage | vin, date, protocol, daily mileage, source mileage | temporary vehicle keys, generic metric key/value rows, per-frame raw details | | MySQL | `vehicle_identity_binding` | Manual identity mapping | vin, plate, phone, device id | registration history | | MySQL | `jt808_registration` | JT808 registration and auth trace | phone, device id, plate, auth code, vin match state, first/latest seen | GB32960 or MQTT records | -| Redis | `vehicle:latest:{vehicleKey}` | Latest merged realtime state | cross-protocol latest fields | historical data | +| Redis | `vehicle:latest:{vehicleKey}` | Latest merged realtime state | cross-protocol latest fields only | full parsed payloads and historical data | | Redis | `vehicle:realtime-raw:{protocol}:{vehicleKey}` | Latest full realtime protocol state | latest protocol parsed payload | historical data | ## Event Envelope Rules diff --git a/docs/architecture/production-data-plane-inventory.md b/docs/architecture/production-data-plane-inventory.md index 074c7249..0a675304 100644 --- a/docs/architecture/production-data-plane-inventory.md +++ b/docs/architecture/production-data-plane-inventory.md @@ -81,9 +81,9 @@ Redis 使用 DB 50,定位为实时缓存,不作为历史事实来源。 | Key 族 | 用途 | | --- | --- | -| `vehicle:latest:{vehicleKey}` | 跨协议合并后的最新实时快照 | +| `vehicle:latest:{vehicleKey}` | 跨协议合并后的最新核心字段快照,不重复保存完整 parsed | | `vehicle:latest:{vehicleKey}:{protocol}` | 单协议最新快照 | -| `vehicle:realtime-raw:{protocol}:{vehicleKey}` | 单协议最新完整 parsed 状态 | +| `vehicle:realtime-raw:{protocol}:{vehicleKey}` | 单协议最新完整 parsed 状态,是实时完整协议字段的唯一 Redis 副本 | | `vehicle:online:{vehicleKey}` | 在线状态和 TTL | | `vehicle:protocols:{vehicleKey}` | 当前车辆最近出现过的协议集合 | | `vehicle:last_seen` | 最近活跃车辆排序集合 | diff --git a/docs/go-version-data-flow.html b/docs/go-version-data-flow.html index 27d738eb..91133eb3 100644 --- a/docs/go-version-data-flow.html +++ b/docs/go-version-data-flow.html @@ -638,7 +638,7 @@ 接口类型数据源说明 VIN 是否在线Redis online keyTTL 内有数据即在线,默认 TTL 600 秒。 - VIN 实时数据Redis merged snapshot跨协议合并 fields,并保留各协议 ProtocolData。 + VIN 实时数据Redis latest snapshot跨协议合并核心 fields;完整协议字段通过 realtime-raw 查询。 单协议实时 RAWRedis realtime-raw查看某 VIN/phone 在某协议下最新 parsed 全量字段。 @@ -722,7 +722,7 @@
实时多协议合并
- Redis 同时保留 latest snapshot、protocol snapshot、realtime-raw,既能看统一实时,也能追单协议原始字段。 + Redis latest snapshot 只保留统一实时核心字段;完整协议 parsed 只放 realtime-raw,避免重复缓存大 JSON。
身份解析降级
diff --git a/go/vehicle-gateway/internal/realtime/model.go b/go/vehicle-gateway/internal/realtime/model.go index 3474f4af..342280d3 100644 --- a/go/vehicle-gateway/internal/realtime/model.go +++ b/go/vehicle-gateway/internal/realtime/model.go @@ -7,19 +7,18 @@ import ( ) type Snapshot struct { - VehicleKey string `json:"vehicle_key"` - VIN string `json:"vin"` - Protocol envelope.Protocol `json:"protocol,omitempty"` - Protocols []envelope.Protocol `json:"protocols,omitempty"` - EventID string `json:"event_id,omitempty"` - EventTimeMS int64 `json:"event_time_ms"` - ReceivedAtMS int64 `json:"received_at_ms"` - SourceEndpoint string `json:"source_endpoint,omitempty"` - Fields map[string]any `json:"fields,omitempty"` - FieldTimesMS map[string]int64 `json:"field_times_ms,omitempty"` - Parsed map[string]any `json:"parsed,omitempty"` - ProtocolData map[envelope.Protocol]map[string]any `json:"protocol_data,omitempty"` - UpdatedAtMS int64 `json:"updated_at_ms"` + VehicleKey string `json:"vehicle_key"` + VIN string `json:"vin"` + Protocol envelope.Protocol `json:"protocol,omitempty"` + Protocols []envelope.Protocol `json:"protocols,omitempty"` + EventID string `json:"event_id,omitempty"` + EventTimeMS int64 `json:"event_time_ms"` + ReceivedAtMS int64 `json:"received_at_ms"` + SourceEndpoint string `json:"source_endpoint,omitempty"` + Fields map[string]any `json:"fields,omitempty"` + FieldTimesMS map[string]int64 `json:"field_times_ms,omitempty"` + Parsed map[string]any `json:"parsed,omitempty"` + UpdatedAtMS int64 `json:"updated_at_ms"` } type OnlineStatus struct { diff --git a/go/vehicle-gateway/internal/realtime/repository.go b/go/vehicle-gateway/internal/realtime/repository.go index db500266..4019026d 100644 --- a/go/vehicle-gateway/internal/realtime/repository.go +++ b/go/vehicle-gateway/internal/realtime/repository.go @@ -92,7 +92,6 @@ func (r *Repository) Update(ctx context.Context, env envelope.FrameEnvelope) err VIN: vin, Fields: map[string]any{}, FieldTimesMS: map[string]int64{}, - ProtocolData: map[envelope.Protocol]map[string]any{}, } } else if merged.VIN == "" && vin != "" { merged.VIN = vin @@ -105,10 +104,6 @@ func (r *Repository) Update(ctx context.Context, env envelope.FrameEnvelope) err merged.SourceEndpoint = env.SourceEndpoint } merged.Protocols = protocols - if merged.ProtocolData == nil { - merged.ProtocolData = map[envelope.Protocol]map[string]any{} - } - merged.ProtocolData[env.Protocol] = cloneMap(protocolSnapshot.Parsed) merged.UpdatedAtMS = nowMS if err := r.setJSON(ctx, mergedKey(vehicleKey), merged, r.cfg.ttl()); err != nil { return err diff --git a/go/vehicle-gateway/internal/realtime/repository_test.go b/go/vehicle-gateway/internal/realtime/repository_test.go index 0e1e4356..b807e94e 100644 --- a/go/vehicle-gateway/internal/realtime/repository_test.go +++ b/go/vehicle-gateway/internal/realtime/repository_test.go @@ -2,6 +2,7 @@ package realtime import ( "context" + "encoding/json" "net/http" "net/http/httptest" "strings" @@ -120,12 +121,12 @@ func TestRepositoryStoresFullParsedProtocolSnapshotAndMergesGB32960Units(t *test if err != nil { t.Fatalf("GetMerged() error = %v", err) } - gb, ok := merged.ProtocolData[envelope.ProtocolGB32960] - if !ok { - t.Fatalf("merged protocol data missing: %#v", merged.ProtocolData) + mergedJSON, err := json.Marshal(merged) + if err != nil { + t.Fatal(err) } - if len(gb["data_units"].([]any)) != 2 { - t.Fatalf("merged protocol data did not keep full parsed fields: %#v", gb) + if strings.Contains(string(mergedJSON), "protocol_data") { + t.Fatalf("merged snapshot should not duplicate full protocol parsed data: %s", string(mergedJSON)) } }