Files
lingniu-vehicle-ingest/go/vehicle-gateway/internal/stats/daily_metric.go
2026-07-02 00:18:35 +08:00

276 lines
7.4 KiB
Go

package stats
import (
"context"
"database/sql"
"errors"
"strconv"
"strings"
"time"
"lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope"
)
const (
MetricDailyMileageKM = "daily_mileage_km"
MetricDailyTotalMileageKM = "daily_total_mileage_km"
)
type Execer interface {
ExecContext(context.Context, string, ...any) (sql.Result, error)
}
type Writer struct {
exec Execer
loc *time.Location
}
type MetricSample struct {
VehicleKey string
VIN string
Protocol envelope.Protocol
StatDate string
MetricKey string
MetricValue float64
TotalMileageKM float64
}
func NewWriter(exec Execer, loc *time.Location) *Writer {
if exec == nil {
panic("stats execer must not be nil")
}
if loc == nil {
loc = time.FixedZone("Asia/Shanghai", 8*3600)
}
return &Writer{exec: exec, loc: loc}
}
func (w *Writer) EnsureSchema(ctx context.Context) error {
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 {
samples, err := SamplesFromEnvelope(env, w.loc)
if err != nil {
return err
}
for _, sample := range samples {
if _, err := w.exec.ExecContext(ctx, upsertDailyMetricSQL,
sample.VehicleKey,
sample.VIN,
sample.StatDate,
string(sample.Protocol),
sample.MetricKey,
sample.MetricValue,
sample.TotalMileageKM,
sample.TotalMileageKM); err != nil {
return err
}
}
return nil
}
func SamplesFromEnvelope(env envelope.FrameEnvelope, loc *time.Location) ([]MetricSample, error) {
vin := strings.TrimSpace(env.VIN)
vehicleKey := strings.TrimSpace(env.VehicleKey())
if vehicleKey == "" || strings.HasSuffix(vehicleKey, ":unknown") {
return nil, nil
}
totalMileage, ok := floatField(env, envelope.FieldTotalMileageKM)
if !ok {
return nil, nil
}
if totalMileage <= 0 {
return nil, nil
}
if loc == nil {
loc = time.FixedZone("Asia/Shanghai", 8*3600)
}
eventMS := env.EventTimeMS
if eventMS <= 0 {
eventMS = env.ReceivedAtMS
}
if eventMS <= 0 {
return nil, errors.New("event or received time is required")
}
statDate := time.UnixMilli(eventMS).In(loc).Format("2006-01-02")
return []MetricSample{
{
VehicleKey: vehicleKey,
VIN: vin,
Protocol: env.Protocol,
StatDate: statDate,
MetricKey: MetricDailyMileageKM,
MetricValue: 0,
TotalMileageKM: totalMileage,
},
{
VehicleKey: vehicleKey,
VIN: vin,
Protocol: env.Protocol,
StatDate: statDate,
MetricKey: MetricDailyTotalMileageKM,
MetricValue: totalMileage,
TotalMileageKM: totalMileage,
},
}, nil
}
const upsertDailyMetricSQL = `
INSERT INTO vehicle_daily_metric
(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')
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)
ELSE LEAST(first_total_mileage_km, VALUES(first_total_mileage_km))
END,
latest_total_mileage_km = CASE
WHEN latest_total_mileage_km IS NULL OR latest_total_mileage_km <= 0
THEN VALUES(latest_total_mileage_km)
ELSE GREATEST(latest_total_mileage_km, VALUES(latest_total_mileage_km))
END,
metric_value = CASE
WHEN metric_key = 'daily_mileage_km'
THEN GREATEST(
CASE
WHEN latest_total_mileage_km IS NULL OR latest_total_mileage_km <= 0
THEN VALUES(latest_total_mileage_km)
ELSE latest_total_mileage_km
END,
VALUES(latest_total_mileage_km)
)
- CASE
WHEN first_total_mileage_km IS NULL OR first_total_mileage_km <= 0
THEN VALUES(first_total_mileage_km)
ELSE LEAST(first_total_mileage_km, VALUES(first_total_mileage_km))
END
WHEN metric_key = 'daily_total_mileage_km'
THEN GREATEST(
CASE
WHEN latest_total_mileage_km IS NULL OR latest_total_mileage_km <= 0
THEN VALUES(latest_total_mileage_km)
ELSE latest_total_mileage_km
END,
VALUES(latest_total_mileage_km)
)
ELSE VALUES(metric_value)
END,
sample_count = sample_count + 1,
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
}
value, ok := env.Fields[key]
if !ok || value == nil {
return 0, false
}
switch typed := value.(type) {
case float64:
return typed, true
case float32:
return float64(typed), true
case int:
return float64(typed), true
case int64:
return float64(typed), true
case uint16:
return float64(typed), true
case uint32:
return float64(typed), true
case string:
parsed, err := strconv.ParseFloat(strings.TrimSpace(typed), 64)
return parsed, err == nil
default:
return 0, false
}
}