Files
lingniu-vehicle-ingest/go/vehicle-gateway/internal/stats/hydrogen_energy_parameter.go

143 lines
5.3 KiB
Go

package stats
import (
"context"
"database/sql"
"fmt"
"sort"
"strings"
)
const HydrogenEnergyParameterTableSQL = `CREATE TABLE IF NOT EXISTS vehicle_hydrogen_energy_parameter (
vin VARCHAR(64) NOT NULL,
battery_capacity_kwh DECIMAL(12,3) NOT NULL,
hydrogen_energy_kwh_per_kg DECIMAL(12,3) NOT NULL DEFAULT 16.000,
confirmed_by VARCHAR(96) NOT NULL DEFAULT '',
source_note VARCHAR(255) NOT NULL DEFAULT '',
effective_from DATE NOT NULL,
effective_to DATE NULL,
active TINYINT(1) NOT NULL DEFAULT 1,
created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
updated_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3),
PRIMARY KEY (vin, effective_from),
KEY idx_hydrogen_energy_parameter_active (active, effective_from, effective_to, vin),
CONSTRAINT chk_hydrogen_energy_parameter_capacity CHECK (battery_capacity_kwh > 0 AND battery_capacity_kwh <= 500),
CONSTRAINT chk_hydrogen_energy_parameter_conversion CHECK (hydrogen_energy_kwh_per_kg > 0 AND hydrogen_energy_kwh_per_kg <= 50),
CONSTRAINT chk_hydrogen_energy_parameter_dates CHECK (effective_to IS NULL OR effective_to >= effective_from)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci`
const (
hydrogenEnergyAssetSourceNote = "ln_asset_management.vehicle_model.battery_capacity"
hydrogenEnergyAssetConfirmedBy = "资产车型主数据自动同步"
hydrogenEnergyAssetEffectiveAt = "1970-01-01"
)
type HydrogenEnergyParameterSyncResult struct {
Read int
Written int
Deactivated int64
}
type hydrogenEnergyParameterRow struct {
VIN string
SourceVehicleID int64
CapacityKWh float64
SourceUpdatedAt sql.NullTime
}
// SyncHydrogenEnergyParameters projects each active hydrogen vehicle's rated
// battery capacity from the asset model master into the VIN-scoped calculation
// parameter table. The fixed, oldest effective date lets a later effective-dated
// business override win without disabling automatic master-data synchronization.
func SyncHydrogenEnergyParameters(ctx context.Context, db *sql.DB, sourceSchema string) (HydrogenEnergyParameterSyncResult, error) {
result := HydrogenEnergyParameterSyncResult{}
if db == nil {
return result, fmt.Errorf("mysql database is required")
}
sourceSchema = strings.TrimSpace(sourceSchema)
if !mysqlIdentifier.MatchString(sourceSchema) {
return result, fmt.Errorf("invalid asset schema %q", sourceSchema)
}
query := fmt.Sprintf(`SELECT vi.id,UPPER(TRIM(vi.vin)),vm.battery_capacity,
CASE
WHEN vi.update_time IS NULL THEN vm.update_time
WHEN vm.update_time IS NULL THEN vi.update_time
ELSE GREATEST(vi.update_time,vm.update_time)
END
FROM %s.vehicle_info vi
JOIN %s.vehicle_model vm ON vm.id=vi.vehicle_model_id
WHERE COALESCE(vi.del_flag,'0')='0' AND COALESCE(vm.del_flag,'0')='0'
AND LENGTH(TRIM(vi.vin))=17
AND vm.tank_capacity>0
AND vm.battery_capacity>0 AND vm.battery_capacity<=500`, sourceSchema, sourceSchema)
rows, err := db.QueryContext(ctx, query)
if err != nil {
return result, fmt.Errorf("read asset hydrogen energy parameters: %w", err)
}
byVIN := map[string]hydrogenEnergyParameterRow{}
for rows.Next() {
var row hydrogenEnergyParameterRow
if err := rows.Scan(&row.SourceVehicleID, &row.VIN, &row.CapacityKWh, &row.SourceUpdatedAt); err != nil {
rows.Close()
return result, err
}
result.Read++
if row.CapacityKWh <= 0 || row.CapacityKWh > 500 {
continue
}
if previous, exists := byVIN[row.VIN]; !exists || newerEnergyParameterRow(row, previous) {
byVIN[row.VIN] = row
}
}
if err := rows.Close(); err != nil {
return result, err
}
tx, err := db.BeginTx(ctx, nil)
if err != nil {
return result, err
}
defer tx.Rollback()
if _, err := tx.ExecContext(ctx, `UPDATE vehicle_hydrogen_energy_parameter
SET active=0
WHERE active=1 AND source_note=?`, hydrogenEnergyAssetSourceNote); err != nil {
return result, err
}
const upsert = `INSERT INTO vehicle_hydrogen_energy_parameter(
vin,battery_capacity_kwh,hydrogen_energy_kwh_per_kg,confirmed_by,source_note,effective_from,effective_to,active
) VALUES(?,?,16.000,?,?,?,NULL,1)
ON DUPLICATE KEY UPDATE battery_capacity_kwh=VALUES(battery_capacity_kwh),
confirmed_by=VALUES(confirmed_by),source_note=VALUES(source_note),effective_to=NULL,active=1`
vins := make([]string, 0, len(byVIN))
for vin := range byVIN {
vins = append(vins, vin)
}
sort.Strings(vins)
for _, vin := range vins {
row := byVIN[vin]
if _, err := tx.ExecContext(ctx, upsert, row.VIN, row.CapacityKWh, hydrogenEnergyAssetConfirmedBy, hydrogenEnergyAssetSourceNote, hydrogenEnergyAssetEffectiveAt); err != nil {
return result, err
}
result.Written++
}
if err := tx.QueryRowContext(ctx, `SELECT COUNT(*)
FROM vehicle_hydrogen_energy_parameter
WHERE active=0 AND source_note=?`, hydrogenEnergyAssetSourceNote).Scan(&result.Deactivated); err != nil {
return result, err
}
if err := tx.Commit(); err != nil {
return result, err
}
return result, nil
}
func newerEnergyParameterRow(left, right hydrogenEnergyParameterRow) bool {
if left.SourceUpdatedAt.Valid != right.SourceUpdatedAt.Valid {
return left.SourceUpdatedAt.Valid
}
if left.SourceUpdatedAt.Valid && !left.SourceUpdatedAt.Time.Equal(right.SourceUpdatedAt.Time) {
return left.SourceUpdatedAt.Time.After(right.SourceUpdatedAt.Time)
}
return left.SourceVehicleID > right.SourceVehicleID
}