From f9b818294954bd478e49bba294f76533098d8a2c Mon Sep 17 00:00:00 2001 From: lingniu Date: Fri, 3 Jul 2026 21:22:21 +0800 Subject: [PATCH] feat(platform): query tdengine history data --- .../apps/api/cmd/platform-api/main.go | 1 + vehicle-data-platform/apps/api/go.mod | 14 +- vehicle-data-platform/apps/api/go.sum | 31 ++++ .../apps/api/internal/app/server.go | 14 +- .../apps/api/internal/config/config.go | 2 + .../apps/api/internal/platform/mock_store.go | 4 + .../api/internal/platform/production_store.go | 147 ++++++++++++++++-- .../internal/platform/query_builders_test.go | 4 +- .../apps/api/internal/platform/service.go | 3 +- .../api/internal/platform/tdengine_queries.go | 138 +++++++++++----- vehicle-data-platform/docs/deployment.md | 3 +- 11 files changed, 306 insertions(+), 55 deletions(-) diff --git a/vehicle-data-platform/apps/api/cmd/platform-api/main.go b/vehicle-data-platform/apps/api/cmd/platform-api/main.go index ac161957..32b1aa9a 100644 --- a/vehicle-data-platform/apps/api/cmd/platform-api/main.go +++ b/vehicle-data-platform/apps/api/cmd/platform-api/main.go @@ -5,6 +5,7 @@ import ( "net/http" _ "github.com/go-sql-driver/mysql" + _ "github.com/taosdata/driver-go/v3/taosWS" "lingniu/vehicle-data-platform/apps/api/internal/app" "lingniu/vehicle-data-platform/apps/api/internal/config" diff --git a/vehicle-data-platform/apps/api/go.mod b/vehicle-data-platform/apps/api/go.mod index 214f6d55..ae8bac84 100644 --- a/vehicle-data-platform/apps/api/go.mod +++ b/vehicle-data-platform/apps/api/go.mod @@ -2,6 +2,16 @@ module lingniu/vehicle-data-platform/apps/api go 1.26 -require github.com/go-sql-driver/mysql v1.9.3 +require ( + github.com/go-sql-driver/mysql v1.9.3 + github.com/taosdata/driver-go/v3 v3.8.1 +) -require filippo.io/edwards25519 v1.1.0 // indirect +require ( + filippo.io/edwards25519 v1.1.0 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/gorilla/websocket v1.5.0 // indirect + github.com/json-iterator/go v1.1.12 // indirect + github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421 // indirect + github.com/modern-go/reflect2 v1.0.2 // indirect +) diff --git a/vehicle-data-platform/apps/api/go.sum b/vehicle-data-platform/apps/api/go.sum index 4bcdcfa5..c3ba046d 100644 --- a/vehicle-data-platform/apps/api/go.sum +++ b/vehicle-data-platform/apps/api/go.sum @@ -1,4 +1,35 @@ filippo.io/edwards25519 v1.1.0 h1:FNf4tywRC1HmFuKW5xopWpigGjJKiJSV0Cqo0cJWDaA= filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/go-sql-driver/mysql v1.9.3 h1:U/N249h2WzJ3Ukj8SowVFjdtZKfu9vlLZxjPXV1aweo= github.com/go-sql-driver/mysql v1.9.3/go.mod h1:qn46aNg1333BRMNU69Lq93t8du/dwxI64Gl8i5p1WMU= +github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/gorilla/websocket v1.5.0 h1:PPwGk2jz7EePpoHN/+ClbZu8SPxiqlu12wZP/3sWmnc= +github.com/gorilla/websocket v1.5.0/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= +github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= +github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= +github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421 h1:ZqeYNhU3OHLH3mGKHDcjJRFFRrJa6eAM5H+CtDdOsPc= +github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= +github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= +github.com/stretchr/objx v0.5.0 h1:1zr/of2m5FGMsad5YfcqgdqdWrIhu+EBEJRhR1U7z/c= +github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.0/go.mod h1:yNjHg4UonilssWZ8iaSj1OCr/vHnekPRkoO+kdMU+MU= +github.com/stretchr/testify v1.8.2 h1:+h33VjcLVPDHtOdpUCuF+7gSuG3yGIftsP1YvFihtJ8= +github.com/stretchr/testify v1.8.2/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= +github.com/taosdata/driver-go/v3 v3.8.1 h1:kkd4ABsGiU+oXDbsw/sic985LKAvnpF9Gb/TEunTnLE= +github.com/taosdata/driver-go/v3 v3.8.1/go.mod h1:S6OGOinfR0xxxaMGsvBi9cLkYxEIW1p6qqr8QJATTlg= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/vehicle-data-platform/apps/api/internal/app/server.go b/vehicle-data-platform/apps/api/internal/app/server.go index e706b490..e04a8541 100644 --- a/vehicle-data-platform/apps/api/internal/app/server.go +++ b/vehicle-data-platform/apps/api/internal/app/server.go @@ -2,6 +2,7 @@ package app import ( "context" + "database/sql" "log" "net/http" "time" @@ -16,11 +17,20 @@ func NewServer(cfg config.Config) http.Handler { if cfg.MySQLDSN != "" { ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() - db, err := platform.OpenMySQL(ctx, cfg.MySQLDSN) + db, err := platform.OpenSQL(ctx, "mysql", cfg.MySQLDSN) if err != nil { log.Printf("production mysql store disabled: %v", err) } else { - store = platform.NewProductionStore(db, cfg.TDengineDatabase) + var tdengine *sql.DB + if cfg.TDengineDSN != "" { + tdengine, err = platform.OpenSQL(ctx, cfg.TDengineDriver, cfg.TDengineDSN) + if err != nil { + log.Printf("production tdengine store disabled: %v", err) + } else { + log.Printf("production tdengine store enabled") + } + } + store = platform.NewProductionStore(db, tdengine, cfg.TDengineDatabase) log.Printf("production mysql store enabled") } } diff --git a/vehicle-data-platform/apps/api/internal/config/config.go b/vehicle-data-platform/apps/api/internal/config/config.go index 834a14aa..8274c46d 100644 --- a/vehicle-data-platform/apps/api/internal/config/config.go +++ b/vehicle-data-platform/apps/api/internal/config/config.go @@ -13,6 +13,7 @@ type Config struct { RedisUsername string RedisPassword string RedisDB int + TDengineDriver string TDengineDSN string TDengineDatabase string CapacityCheckBin string @@ -28,6 +29,7 @@ func Load() Config { RedisUsername: os.Getenv("REDIS_USERNAME"), RedisPassword: os.Getenv("REDIS_PASSWORD"), RedisDB: envInt("REDIS_DB", 50), + TDengineDriver: env("TDENGINE_DRIVER", "taosWS"), TDengineDSN: os.Getenv("TDENGINE_DSN"), TDengineDatabase: env("TDENGINE_DATABASE", "lingniu_vehicle_ts"), CapacityCheckBin: os.Getenv("CAPACITY_CHECK_BIN"), diff --git a/vehicle-data-platform/apps/api/internal/platform/mock_store.go b/vehicle-data-platform/apps/api/internal/platform/mock_store.go index 877a5779..c1338e29 100644 --- a/vehicle-data-platform/apps/api/internal/platform/mock_store.go +++ b/vehicle-data-platform/apps/api/internal/platform/mock_store.go @@ -67,6 +67,10 @@ func (m *MockStore) RealtimeLocations(_ context.Context, query url.Values) (Page } func (m *MockStore) HistoryLocations(ctx context.Context, query url.Values) (Page[HistoryLocationRow], error) { + return m.HistoryLocationsFromTDengine(ctx, query) +} + +func (m *MockStore) HistoryLocationsFromTDengine(ctx context.Context, query url.Values) (Page[HistoryLocationRow], error) { realtime, _ := m.RealtimeLocations(ctx, query) rows := make([]HistoryLocationRow, 0, len(realtime.Items)) for _, row := range realtime.Items { diff --git a/vehicle-data-platform/apps/api/internal/platform/production_store.go b/vehicle-data-platform/apps/api/internal/platform/production_store.go index 684cbc1d..14e0333d 100644 --- a/vehicle-data-platform/apps/api/internal/platform/production_store.go +++ b/vehicle-data-platform/apps/api/internal/platform/production_store.go @@ -6,24 +6,26 @@ import ( "encoding/json" "fmt" "net/url" + "strconv" "strings" "time" ) type ProductionStore struct { - db *sql.DB - database string + db *sql.DB + tdengine *sql.DB + tdDatabase string } -func NewProductionStore(db *sql.DB, tdengineDatabase string) *ProductionStore { +func NewProductionStore(db *sql.DB, tdengine *sql.DB, tdengineDatabase string) *ProductionStore { if db == nil { panic("production db must not be nil") } - return &ProductionStore{db: db, database: tdengineDatabase} + return &ProductionStore{db: db, tdengine: tdengine, tdDatabase: tdengineDatabase} } -func OpenMySQL(ctx context.Context, dsn string) (*sql.DB, error) { - db, err := sql.Open("mysql", dsn) +func OpenSQL(ctx context.Context, driver, dsn string) (*sql.DB, error) { + db, err := sql.Open(driver, dsn) if err != nil { return nil, err } @@ -139,10 +141,86 @@ func (s *ProductionStore) HistoryLocations(ctx context.Context, query url.Values return Page[HistoryLocationRow]{Items: items, Total: realtime.Total, Limit: realtime.Limit, Offset: realtime.Offset}, nil } +func (s *ProductionStore) HistoryLocationsFromTDengine(ctx context.Context, query url.Values) (Page[HistoryLocationRow], error) { + if s.tdengine == nil { + return s.HistoryLocations(ctx, query) + } + limit, offset := buildLimitOffset(query) + tdQuery := map[string]string{ + "protocol": query.Get("protocol"), + "vin": query.Get("vin"), + "limit": strconv.Itoa(limit), + "offset": strconv.Itoa(offset), + } + if value := strings.TrimSpace(query.Get("dateFrom")); value != "" { + tdQuery["dateFrom"] = value + } + if value := strings.TrimSpace(query.Get("dateTo")); value != "" { + tdQuery["dateTo"] = value + } + built := buildHistoryLocationSQL(s.tdDatabase, tdQuery) + rows, err := s.tdengine.QueryContext(ctx, built.Text, built.Args...) + if err != nil { + return Page[HistoryLocationRow]{}, err + } + defer rows.Close() + items := make([]HistoryLocationRow, 0) + for rows.Next() { + var row HistoryLocationRow + var ts, receivedAt string + var longitude, latitude, speed, mileage sql.NullFloat64 + if err := rows.Scan(&ts, &row.VIN, &row.Protocol, &longitude, &latitude, &speed, &mileage, &receivedAt); err != nil { + return Page[HistoryLocationRow]{}, err + } + row.Longitude = nullFloat64(longitude) + row.Latitude = nullFloat64(latitude) + row.SpeedKmh = nullFloat64(speed) + row.TotalMileageKm = nullFloat64(mileage) + row.DeviceTime = ts + row.ServerTime = firstNonEmpty(receivedAt, ts) + items = append(items, row) + } + if err := rows.Err(); err != nil { + return Page[HistoryLocationRow]{}, err + } + return Page[HistoryLocationRow]{Items: items, Total: len(items), Limit: limit, Offset: offset}, nil +} + func (s *ProductionStore) RawFrames(ctx context.Context, query RawFrameQuery) (Page[RawFrameRow], error) { - _ = ctx - _ = buildRawFrameSQL(s.database, query) - return Page[RawFrameRow]{Items: []RawFrameRow{}, Total: 0, Limit: query.Limit, Offset: query.Offset}, nil + if s.tdengine == nil { + return Page[RawFrameRow]{Items: []RawFrameRow{}, Total: 0, Limit: query.Limit, Offset: query.Offset}, nil + } + built := buildRawFrameSQL(s.tdDatabase, query) + rows, err := s.tdengine.QueryContext(ctx, built.Text, built.Args...) + if err != nil { + return Page[RawFrameRow]{}, err + } + defer rows.Close() + items := make([]RawFrameRow, 0) + for rows.Next() { + var row RawFrameRow + var ts, frameID, eventTime, receivedAt, parsedFields, parseStatus, parseError, sourceEndpoint, protocol, vehicleKey, vin, phone string + var rawSizeBytes int + if err := rows.Scan(&ts, &frameID, &eventTime, &receivedAt, &rawSizeBytes, &parsedFields, &parseStatus, &parseError, &sourceEndpoint, &protocol, &vehicleKey, &vin, &phone); err != nil { + return Page[RawFrameRow]{}, err + } + row.ID = frameID + row.VIN = vin + row.Protocol = protocol + row.FrameType = vehicleKey + row.DeviceTime = firstNonEmpty(eventTime, ts) + row.ServerTime = firstNonEmpty(receivedAt, ts) + row.RawSizeBytes = rawSizeBytes + row.ParsedFields = parsedFieldsFromString(parsedFields) + if len(query.Fields) > 0 { + row.ParsedFields = filterParsedFieldsMap(row.ParsedFields, query.Fields) + } + items = append(items, row) + } + if err := rows.Err(); err != nil { + return Page[RawFrameRow]{}, err + } + return Page[RawFrameRow]{Items: items, Total: len(items), Limit: query.Limit, Offset: query.Offset}, nil } func (s *ProductionStore) DailyMileage(ctx context.Context, query url.Values) (Page[DailyMileageRow], error) { @@ -206,11 +284,11 @@ func (s *ProductionStore) OpsHealth(ctx context.Context) (OpsHealth, error) { {Name: "MySQL realtime", Status: mysqlStatus, Detail: mysqlDetail}, {Name: "vehicle_realtime_snapshot", Status: "ok", Detail: "读取实时快照表"}, {Name: "vehicle_realtime_location", Status: "ok", Detail: "读取实时位置表"}, - {Name: "TDengine raw_frames", Status: "warning", Detail: "中台 RAW 查询接口待接 TDengine 驱动"}, + {Name: "TDengine raw_frames", Status: tdengineStatus(s.tdengine), Detail: tdengineDetail(s.tdengine)}, }, KafkaLag: 0, RedisOnlineKeys: redisKeys, - TDengineWritable: false, + TDengineWritable: s.tdengine != nil, MySQLWritable: mysqlStatus == "ok", }, nil } @@ -242,3 +320,50 @@ func parsedFieldsFromString(value string) map[string]any { } return fields } + +func filterParsedFieldsMap(fields map[string]any, names []string) map[string]any { + if len(names) == 0 || len(fields) == 0 { + return fields + } + out := make(map[string]any, len(names)) + for _, name := range names { + name = strings.TrimSpace(name) + if name == "" { + continue + } + if value, ok := fields[name]; ok { + out[name] = value + } + } + return out +} + +func nullFloat64(value sql.NullFloat64) float64 { + if value.Valid { + return value.Float64 + } + return 0 +} + +func firstNonEmpty(values ...string) string { + for _, value := range values { + if strings.TrimSpace(value) != "" { + return value + } + } + return "" +} + +func tdengineStatus(db *sql.DB) string { + if db == nil { + return "warning" + } + return "ok" +} + +func tdengineDetail(db *sql.DB) string { + if db == nil { + return "TDengine 未配置或连接失败,历史查询降级" + } + return "读取 TDengine raw_frames / vehicle_locations" +} diff --git a/vehicle-data-platform/apps/api/internal/platform/query_builders_test.go b/vehicle-data-platform/apps/api/internal/platform/query_builders_test.go index 6f72ec2f..882ae2d4 100644 --- a/vehicle-data-platform/apps/api/internal/platform/query_builders_test.go +++ b/vehicle-data-platform/apps/api/internal/platform/query_builders_test.go @@ -36,10 +36,10 @@ func TestBuildRawFrameSQL(t *testing.T) { VIN: "VIN001", Limit: 1, }) - if !strings.Contains(built.Text, "lingniu_vehicle_ts.raw_frames") || !strings.Contains(built.Text, "parsed_fields") { + if !strings.Contains(built.Text, "lingniu_vehicle_ts.raw_gb32960_") || !strings.Contains(built.Text, "parsed_fields") { t.Fatalf("SQL = %s", built.Text) } - if len(built.Args) != 4 || built.Args[0] != "GB32960" || built.Args[1] != "VIN001" || built.Args[2] != 1 { + if len(built.Args) != 0 || !strings.Contains(built.Text, "protocol = 'GB32960'") || !strings.Contains(built.Text, "vin = 'VIN001'") || !strings.Contains(built.Text, "LIMIT 1 OFFSET 0") { t.Fatalf("args = %#v", built.Args) } } diff --git a/vehicle-data-platform/apps/api/internal/platform/service.go b/vehicle-data-platform/apps/api/internal/platform/service.go index 9bdf989a..42b0d6c7 100644 --- a/vehicle-data-platform/apps/api/internal/platform/service.go +++ b/vehicle-data-platform/apps/api/internal/platform/service.go @@ -10,6 +10,7 @@ type Store interface { Vehicles(context.Context, url.Values) (Page[VehicleRow], error) RealtimeLocations(context.Context, url.Values) (Page[RealtimeLocationRow], error) HistoryLocations(context.Context, url.Values) (Page[HistoryLocationRow], error) + HistoryLocationsFromTDengine(context.Context, url.Values) (Page[HistoryLocationRow], error) RawFrames(context.Context, RawFrameQuery) (Page[RawFrameRow], error) DailyMileage(context.Context, url.Values) (Page[DailyMileageRow], error) QualityIssues(context.Context, url.Values) (Page[QualityIssueRow], error) @@ -48,7 +49,7 @@ func (s *Service) RealtimeLocations(ctx context.Context, query url.Values) (Page } func (s *Service) HistoryLocations(ctx context.Context, query url.Values) (Page[HistoryLocationRow], error) { - return s.store.HistoryLocations(ctx, query) + return s.store.HistoryLocationsFromTDengine(ctx, query) } func (s *Service) RawFrames(ctx context.Context, query RawFrameQuery) (Page[RawFrameRow], error) { diff --git a/vehicle-data-platform/apps/api/internal/platform/tdengine_queries.go b/vehicle-data-platform/apps/api/internal/platform/tdengine_queries.go index 36890e3c..57053714 100644 --- a/vehicle-data-platform/apps/api/internal/platform/tdengine_queries.go +++ b/vehicle-data-platform/apps/api/internal/platform/tdengine_queries.go @@ -1,59 +1,78 @@ package platform -import "strings" +import ( + "crypto/sha1" + "encoding/hex" + "strconv" + "strings" + "time" +) func buildRawFrameSQL(database string, query RawFrameQuery) SQLQuery { - table := qualifyTDengine(database, "raw_frames") - where := []string{"1 = 1"} + table := rawFrameTable(database, query) + where := rawFrameWhere(query) + parsedFieldsSelect := "'' AS parsed_fields" + if query.IncludeFields || len(query.Fields) > 0 { + parsedFieldsSelect = "parsed_json AS parsed_fields" + } args := []any{} - if query.Protocol != "" { - where = append(where, "protocol = ?") - args = append(args, query.Protocol) - } - if query.VIN != "" { - where = append(where, "vin = ?") - args = append(args, query.VIN) - } - if query.DateFrom != "" { - where = append(where, "ts >= ?") - args = append(args, query.DateFrom) - } - if query.DateTo != "" { - where = append(where, "ts <= ?") - args = append(args, query.DateTo) - } limit := query.Limit if limit <= 0 { limit = 100 } - args = append(args, limit, query.Offset) - return SQLQuery{ - Text: `SELECT ts, frame_id, event_time, received_at, raw_size_bytes, parsed_fields, parse_status, ` + - `parse_error, source_endpoint, protocol, vehicle_key, vin, phone FROM ` + table + - ` WHERE ` + strings.Join(where, " AND ") + ` ORDER BY ts DESC LIMIT ? OFFSET ?`, - Args: args, + offset := query.Offset + if offset < 0 { + offset = 0 } + text := `SELECT ts, frame_id, event_time, received_at, raw_size_bytes, ` + parsedFieldsSelect + + `, parse_status, parse_error, source_endpoint, protocol, vehicle_key, vin, phone FROM ` + table + if len(where) > 0 { + text += ` WHERE ` + strings.Join(where, " AND ") + } + text += ` ORDER BY ts DESC LIMIT ` + strconv.Itoa(limit) + ` OFFSET ` + strconv.Itoa(offset) + return SQLQuery{Text: text, Args: args} +} + +func rawFrameWhere(query RawFrameQuery) []string { + where := make([]string, 0, 4) + if query.Protocol != "" { + where = append(where, "protocol = '"+quoteTDengine(strings.ToUpper(strings.TrimSpace(query.Protocol)))+"'") + } + if query.VIN != "" { + where = append(where, "vin = '"+quoteTDengine(strings.TrimSpace(query.VIN))+"'") + } + if query.DateFrom != "" { + where = append(where, "ts >= '"+quoteTDengine(normalizeTDengineTime(query.DateFrom))+"'") + } + if query.DateTo != "" { + where = append(where, "ts <= '"+quoteTDengine(normalizeTDengineTime(query.DateTo))+"'") + } + return where } func buildHistoryLocationSQL(database string, query map[string]string) SQLQuery { - table := qualifyTDengine(database, "vehicle_locations") - where := []string{"1 = 1"} + table := locationTable(database, query) + where := make([]string, 0, 4) args := []any{} if protocol := strings.TrimSpace(query["protocol"]); protocol != "" { - where = append(where, "protocol = ?") - args = append(args, protocol) + where = append(where, "protocol = '"+quoteTDengine(strings.ToUpper(protocol))+"'") } if vin := strings.TrimSpace(query["vin"]); vin != "" { - where = append(where, "vin = ?") - args = append(args, vin) + where = append(where, "vin = '"+quoteTDengine(vin)+"'") + } + if dateFrom := strings.TrimSpace(query["dateFrom"]); dateFrom != "" { + where = append(where, "ts >= '"+quoteTDengine(normalizeTDengineTime(dateFrom))+"'") + } + if dateTo := strings.TrimSpace(query["dateTo"]); dateTo != "" { + where = append(where, "ts <= '"+quoteTDengine(normalizeTDengineTime(dateTo))+"'") } limit, offset := parseLimitOffset(query["limit"], query["offset"]) - args = append(args, limit, offset) - return SQLQuery{ - Text: `SELECT ts, vin, protocol, longitude, latitude, speed_kmh, total_mileage_km, event_time, received_at FROM ` + - table + ` WHERE ` + strings.Join(where, " AND ") + ` ORDER BY ts DESC LIMIT ? OFFSET ?`, - Args: args, + text := `SELECT ts, vin, protocol, longitude, latitude, speed_kmh, total_mileage_km, received_at FROM ` + table + if len(where) > 0 { + text += ` WHERE ` + strings.Join(where, " AND ") } + text += ` ORDER BY ts DESC LIMIT ` + strconv.Itoa(limit) + ` OFFSET ` + strconv.Itoa(offset) + return SQLQuery{Text: text, Args: args} } func qualifyTDengine(database, table string) string { @@ -62,3 +81,50 @@ func qualifyTDengine(database, table string) string { } return database + "." + table } + +func rawFrameTable(database string, query RawFrameQuery) string { + protocol := strings.ToUpper(strings.TrimSpace(query.Protocol)) + vin := strings.TrimSpace(query.VIN) + if protocol == "" || vin == "" { + return qualifyTDengine(database, "raw_frames") + } + if protocol == "JT808" { + return qualifyTDengine(database, "raw_frames") + } + return qualifyTDengine(database, "raw_"+strings.ToLower(protocol)+"_"+hash16(vin)) +} + +func locationTable(database string, query map[string]string) string { + protocol := strings.ToUpper(strings.TrimSpace(query["protocol"])) + vin := strings.TrimSpace(query["vin"]) + if protocol == "" || vin == "" { + return qualifyTDengine(database, "vehicle_locations") + } + return qualifyTDengine(database, "loc_"+strings.ToLower(protocol)+"_"+hash16(vin)) +} + +func hash16(value string) string { + sum := sha1.Sum([]byte(value)) + return hex.EncodeToString(sum[:8]) +} + +func quoteTDengine(value string) string { + return strings.ReplaceAll(value, "'", "''") +} + +func normalizeTDengineTime(value string) string { + value = strings.TrimSpace(value) + if value == "" { + return "" + } + shanghai := time.FixedZone("Asia/Shanghai", 8*3600) + for _, layout := range []string{"2006-01-02T15:04:05", "2006-01-02 15:04:05", "2006-01-02"} { + if parsed, err := time.ParseInLocation(layout, value, shanghai); err == nil { + return parsed.UTC().Format("2006-01-02 15:04:05") + } + } + if parsed, err := time.Parse(time.RFC3339, value); err == nil { + return parsed.UTC().Format("2006-01-02 15:04:05") + } + return value +} diff --git a/vehicle-data-platform/docs/deployment.md b/vehicle-data-platform/docs/deployment.md index 413c60ed..0138bea6 100644 --- a/vehicle-data-platform/docs/deployment.md +++ b/vehicle-data-platform/docs/deployment.md @@ -30,7 +30,8 @@ REDIS_ADDR=r-bp1u741kij7e51i481.redis.rds.aliyuncs.com:6379 REDIS_USERNAME=lingniu_vehicle REDIS_PASSWORD=*** REDIS_DB=50 -TDENGINE_DSN=http://root:***@172.17.111.57:6041 +TDENGINE_DRIVER=taosWS +TDENGINE_DSN=root:***@ws(172.17.111.57:6041)/ TDENGINE_DATABASE=lingniu_vehicle_ts CAPACITY_CHECK_BIN=/opt/lingniu-go-native/current/capacity-check AUTH_TOKEN=***