fix: key daily metrics by vehicle identity

This commit is contained in:
lingniu
2026-07-02 00:18:35 +08:00
parent 9045b871d6
commit bcf7cdae2e
5 changed files with 201 additions and 31 deletions

View File

@@ -26,6 +26,7 @@ type Writer struct {
}
type MetricSample struct {
VehicleKey string
VIN string
Protocol envelope.Protocol
StatDate string
@@ -45,8 +46,13 @@ func NewWriter(exec Execer, loc *time.Location) *Writer {
}
func (w *Writer) EnsureSchema(ctx context.Context) error {
_, err := w.exec.ExecContext(ctx, DailyMetricTableSQL)
return err
if _, err := w.exec.ExecContext(ctx, DailyMetricTableSQL); err != nil {
return err
}
if db, ok := w.exec.(metadataDB); ok {
return migrateDailyMetricTable(ctx, db)
}
return nil
}
func (w *Writer) Append(ctx context.Context, env envelope.FrameEnvelope) error {
@@ -56,6 +62,7 @@ func (w *Writer) Append(ctx context.Context, env envelope.FrameEnvelope) error {
}
for _, sample := range samples {
if _, err := w.exec.ExecContext(ctx, upsertDailyMetricSQL,
sample.VehicleKey,
sample.VIN,
sample.StatDate,
string(sample.Protocol),
@@ -71,7 +78,8 @@ func (w *Writer) Append(ctx context.Context, env envelope.FrameEnvelope) error {
func SamplesFromEnvelope(env envelope.FrameEnvelope, loc *time.Location) ([]MetricSample, error) {
vin := strings.TrimSpace(env.VIN)
if vin == "" {
vehicleKey := strings.TrimSpace(env.VehicleKey())
if vehicleKey == "" || strings.HasSuffix(vehicleKey, ":unknown") {
return nil, nil
}
totalMileage, ok := floatField(env, envelope.FieldTotalMileageKM)
@@ -94,6 +102,7 @@ func SamplesFromEnvelope(env envelope.FrameEnvelope, loc *time.Location) ([]Metr
statDate := time.UnixMilli(eventMS).In(loc).Format("2006-01-02")
return []MetricSample{
{
VehicleKey: vehicleKey,
VIN: vin,
Protocol: env.Protocol,
StatDate: statDate,
@@ -102,6 +111,7 @@ func SamplesFromEnvelope(env envelope.FrameEnvelope, loc *time.Location) ([]Metr
TotalMileageKM: totalMileage,
},
{
VehicleKey: vehicleKey,
VIN: vin,
Protocol: env.Protocol,
StatDate: statDate,
@@ -114,10 +124,11 @@ func SamplesFromEnvelope(env envelope.FrameEnvelope, loc *time.Location) ([]Metr
const upsertDailyMetricSQL = `
INSERT INTO vehicle_daily_metric
(vin, stat_date, protocol, metric_key, metric_value, metric_unit,
(vehicle_key, vin, stat_date, protocol, metric_key, metric_value, metric_unit,
first_total_mileage_km, latest_total_mileage_km, sample_count, calculation_method)
VALUES (?, ?, ?, ?, ?, 'km', ?, ?, 1, 'TOTAL_MILEAGE_DIFF')
VALUES (?, ?, ?, ?, ?, ?, 'km', ?, ?, 1, 'TOTAL_MILEAGE_DIFF')
ON DUPLICATE KEY UPDATE
vin = IF(VALUES(vin) <> '', VALUES(vin), vin),
first_total_mileage_km = CASE
WHEN first_total_mileage_km IS NULL OR first_total_mileage_km <= 0
THEN VALUES(first_total_mileage_km)
@@ -158,6 +169,82 @@ ON DUPLICATE KEY UPDATE
updated_at = CURRENT_TIMESTAMP
`
type metadataDB interface {
Execer
QueryRowContext(context.Context, string, ...any) *sql.Row
}
func migrateDailyMetricTable(ctx context.Context, db metadataDB) error {
columnExists, err := informationSchemaExists(ctx, db, `
SELECT COUNT(*)
FROM information_schema.columns
WHERE table_schema = DATABASE()
AND table_name = 'vehicle_daily_metric'
AND column_name = 'vehicle_key'`)
if err != nil {
return err
}
if !columnExists {
if _, err := db.ExecContext(ctx, `ALTER TABLE vehicle_daily_metric ADD COLUMN vehicle_key VARCHAR(96) NOT NULL DEFAULT '' AFTER id`); err != nil {
return err
}
}
if _, err := db.ExecContext(ctx, `UPDATE vehicle_daily_metric SET vehicle_key = vin WHERE vehicle_key = ''`); err != nil {
return err
}
vehicleIndexExists, err := informationSchemaExists(ctx, db, `
SELECT COUNT(*)
FROM information_schema.statistics
WHERE table_schema = DATABASE()
AND table_name = 'vehicle_daily_metric'
AND index_name = 'uk_daily_metric_vehicle'`)
if err != nil {
return err
}
if !vehicleIndexExists {
if _, err := db.ExecContext(ctx, `ALTER TABLE vehicle_daily_metric ADD UNIQUE KEY uk_daily_metric_vehicle (vehicle_key, stat_date, protocol, metric_key)`); err != nil {
return err
}
}
oldIndexExists, err := informationSchemaExists(ctx, db, `
SELECT COUNT(*)
FROM information_schema.statistics
WHERE table_schema = DATABASE()
AND table_name = 'vehicle_daily_metric'
AND index_name = 'uk_daily_metric'`)
if err != nil {
return err
}
if oldIndexExists {
if _, err := db.ExecContext(ctx, `ALTER TABLE vehicle_daily_metric DROP INDEX uk_daily_metric`); err != nil {
return err
}
}
vinIndexExists, err := informationSchemaExists(ctx, db, `
SELECT COUNT(*)
FROM information_schema.statistics
WHERE table_schema = DATABASE()
AND table_name = 'vehicle_daily_metric'
AND index_name = 'idx_vin'`)
if err != nil {
return err
}
if !vinIndexExists {
if _, err := db.ExecContext(ctx, `ALTER TABLE vehicle_daily_metric ADD KEY idx_vin (vin)`); err != nil {
return err
}
}
return nil
}
func informationSchemaExists(ctx context.Context, db metadataDB, query string) (bool, error) {
var count int
if err := db.QueryRowContext(ctx, query).Scan(&count); err != nil {
return false, err
}
return count > 0, nil
}
func floatField(env envelope.FrameEnvelope, key string) (float64, bool) {
if env.Fields == nil {
return 0, false