From 09f3e688919e6d08b6cfc76cb34d9eba49ad0a2e Mon Sep 17 00:00:00 2001 From: lingniu Date: Fri, 3 Jul 2026 18:37:38 +0800 Subject: [PATCH] feat(go): batch tdengine history writes --- .../tdengine-batch-writer-design.md | 32 ++++- docs/ops/100k-capacity-baseline.md | 3 + docs/ops/go-vehicle-ingest-memory.md | 19 +++ docs/ops/vehicle-ingest-runbook.md | 3 + go/vehicle-gateway/cmd/history-writer/main.go | 126 +++++++++++++++- .../cmd/history-writer/main_test.go | 136 +++++++++++++++++- go/vehicle-gateway/internal/history/writer.go | 93 ++++++++++++ .../internal/history/writer_test.go | 47 ++++++ 8 files changed, 447 insertions(+), 12 deletions(-) diff --git a/docs/architecture/tdengine-batch-writer-design.md b/docs/architecture/tdengine-batch-writer-design.md index 4b79329c..1e199a45 100644 --- a/docs/architecture/tdengine-batch-writer-design.md +++ b/docs/architecture/tdengine-batch-writer-design.md @@ -50,10 +50,9 @@ | 指标 | 含义 | | --- | --- | -| `vehicle_history_batch_flush_total{table,status}` | 批写成功/失败次数 | -| `vehicle_history_batch_rows_total{table,status}` | 批写行数 | -| `vehicle_history_batch_flush_duration_ms{table,status}` | 批写耗时 | -| `vehicle_history_batch_pending_rows{table}` | 当前待 flush 行数 | +| `vehicle_history_batch_flush_total{status}` | 批写成功/失败次数 | +| `vehicle_history_batch_rows_total{status}` | 批写行数 | +| `vehicle_history_batch_flush_duration_ms{status}` | 批写耗时 | | `vehicle_history_kafka_lag` | 下游是否追得上 Kafka | ## Rollout @@ -63,3 +62,28 @@ 3. 在测试环境打开 batch writer。 4. 生产先用小批次 `100/100ms`,观察 TDengine latency 和 Kafka lag。 5. 逐步提升到 `200-500/100ms`。 + +## Implementation Status + +2026-07-03 已完成第一阶段实现: + +- `history.Writer.AppendAllBatch` 按 raw child table 和 location child table 生成多行 `INSERT ... VALUES (...),(...)`。 +- oversized payload chunks 仍写入 `raw_frame_payload_chunks`,同样按 child table 批量写。 +- `history-writer` Kafka 消费默认按 `HISTORY_BATCH_SIZE=200`、`HISTORY_BATCH_WAIT_MS=100` 收集消息。 +- TDengine batch 成功后才批量提交 Kafka messages;batch 失败不提交 offset,让 Kafka 保留可重放语义。 +- 保留 `processHistoryMessage` 和 `AppendAll` 单条路径,便于回退和测试。 + +新增指标: + +```text +vehicle_history_batch_flush_total{status} +vehicle_history_batch_rows_total{status} +vehicle_history_batch_flush_duration_ms{status} +``` + +生产 rollout 建议: + +1. 先部署默认 `200/100ms`。 +2. 观察 `vehicle_history_batch_*` 和 `vehicle_history_kafka_lag`。 +3. 如果 TDengine latency 上升或 Kafka lag 不下降,将 `HISTORY_BATCH_SIZE` 降到 `100`。 +4. 如果稳定且有 backlog,再逐步提升到 `500`。 diff --git a/docs/ops/100k-capacity-baseline.md b/docs/ops/100k-capacity-baseline.md index ff552ec6..a7eb5156 100644 --- a/docs/ops/100k-capacity-baseline.md +++ b/docs/ops/100k-capacity-baseline.md @@ -66,6 +66,9 @@ systemd gateway 已配置 `LimitNOFILE=1048576`,需要持续保持。 | Gateway async sink | `vehicle_async_sink_publish_total{sink,kind,status}` | `error` 不增长 | | Gateway | `vehicle_gateway_publish_total{status="ok"}` | 持续增长 | | Gateway | parse/publish duration | p99 小于容量目标 | +| History writer | `vehicle_history_batch_flush_total{status}` | `error` 不增长 | +| History writer | `vehicle_history_batch_rows_total{status}` | batch 行数持续增长 | +| History writer | `vehicle_history_batch_flush_duration_ms{status}` | flush 延迟不持续上升 | | NATS bridge | `vehicle_bridge_nats_consumer_ack_pending` | 稳态为 0 | | NATS bridge | `vehicle_bridge_nats_consumer_pending` | burst 后下降 | | Kafka consumers | `vehicle_*_kafka_lag` | 稳态为 0,burst 后下降 | diff --git a/docs/ops/go-vehicle-ingest-memory.md b/docs/ops/go-vehicle-ingest-memory.md index ef57f18b..79a86bd3 100644 --- a/docs/ops/go-vehicle-ingest-memory.md +++ b/docs/ops/go-vehicle-ingest-memory.md @@ -406,6 +406,25 @@ vehicle_async_sink_queue_depth{sink} 这些指标是进入带帧率压测前的关键保护栏:如果 `queue_depth` 持续增长或 `enqueue_total{status="timeout"}` 增长,说明入口 publish 队列已成为瓶颈。 +### History writer TDengine 批写 + +2026-07-03 已实现 history-writer 第一阶段批写: + +- `history.Writer.AppendAllBatch` 按 TDengine child table 聚合 raw/location 多行 `INSERT`。 +- `cmd/history-writer` 默认 `HISTORY_BATCH_SIZE=200`、`HISTORY_BATCH_WAIT_MS=100`。 +- TDengine batch 成功后才批量提交 Kafka messages;失败不提交 offset。 +- 保留单条 `AppendAll` 路径作为回退。 + +新增指标: + +```text +vehicle_history_batch_flush_total{status} +vehicle_history_batch_rows_total{status} +vehicle_history_batch_flush_duration_ms{status} +``` + +进入带帧率压测时,必须同时观察 `vehicle_history_kafka_lag` 和 `vehicle_history_batch_*`。 + ## 下次恢复上下文时先做 1. 读本文件。 diff --git a/docs/ops/vehicle-ingest-runbook.md b/docs/ops/vehicle-ingest-runbook.md index f71317f1..96db5240 100644 --- a/docs/ops/vehicle-ingest-runbook.md +++ b/docs/ops/vehicle-ingest-runbook.md @@ -73,6 +73,7 @@ curl -fsS http://127.0.0.1:20214/metrics \ | grep -E 'vehicle_bridge_(nats_consumer|kafka_writes_total|nats_acks_total)' curl -fsS http://127.0.0.1:20212/metrics | grep vehicle_history_kafka_lag +curl -fsS http://127.0.0.1:20212/metrics | grep vehicle_history_batch curl -fsS http://127.0.0.1:20213/metrics | grep vehicle_stat_kafka_lag curl -fsS http://127.0.0.1:20200/metrics | grep vehicle_realtime_kafka_lag ``` @@ -118,6 +119,8 @@ go run ./cmd/load-sim \ | `vehicle_async_sink_queue_depth{sink="nats"}` | 持续增长且不回落 | Gateway 到 NATS/Kafka 的异步 publish 队列开始积压。 | | `vehicle_async_sink_enqueue_total{status="timeout"}` | 任意增长 | Gateway publish 队列已满或 worker 长时间阻塞,入口可能开始丢实时性。 | | `vehicle_async_sink_publish_total{status="error"}` | 连续增长 | NATS/Kafka publish 失败,需要先查中间件连接和日志。 | +| `vehicle_history_batch_flush_total{status="error"}` | 任意增长 | TDengine 批写失败;Kafka offset 不会提交,应先查 TDengine 和 SQL 错误。 | +| `vehicle_history_batch_flush_duration_ms{status="ok"}` | 持续上升 | TDengine 写入延迟增加,可能需要降低 batch size 或扩容 TDengine。 | | `vehicle_bridge_nats_consumer_ack_pending` | 连续 2 分钟 `> 0` | 消息已投递给 bridge,但 Kafka 写入后未完成 ack。 | | `vehicle_bridge_nats_consumer_pending` | 持续增长且 `> 10000` | bridge 消费 NATS 的速度跟不上生产速度。 | | Kafka lag | 连续 5 分钟增长或 `> 10000` | 下游 consumer 或存储存在瓶颈。 | diff --git a/go/vehicle-gateway/cmd/history-writer/main.go b/go/vehicle-gateway/cmd/history-writer/main.go index 4c5f11ba..35dc7a8e 100644 --- a/go/vehicle-gateway/cmd/history-writer/main.go +++ b/go/vehicle-gateway/cmd/history-writer/main.go @@ -6,6 +6,7 @@ import ( "encoding/json" "os" "os/signal" + "strconv" "strings" "syscall" "time" @@ -62,7 +63,9 @@ func main() { logger.Info("history writer started", "driver", cfg.TDengineDriver, "group", cfg.KafkaGroup, - "topics", strings.Join(cfg.KafkaTopics, ",")) + "topics", strings.Join(cfg.KafkaTopics, ","), + "batch_size", cfg.BatchSize, + "batch_wait_ms", cfg.BatchWait) for { message, err := reader.FetchMessage(ctx) @@ -73,7 +76,8 @@ func main() { logger.Error("kafka fetch failed", "error", err) continue } - processHistoryMessage(ctx, logger, registry, writer, reader, message) + batch := collectHistoryBatch(ctx, reader, message, cfg.BatchSize, time.Duration(cfg.BatchWait)*time.Millisecond) + processHistoryBatch(ctx, logger, registry, writer, reader, batch) } } @@ -81,12 +85,42 @@ const kafkaMessageOperationTimeout = 30 * time.Second type historyAppender interface { AppendAll(context.Context, envelope.FrameEnvelope) error + AppendAllBatch(context.Context, []envelope.FrameEnvelope) error } type kafkaMessageCommitter interface { CommitMessages(context.Context, ...kafka.Message) error } +type kafkaMessageFetcher interface { + FetchMessage(context.Context) (kafka.Message, error) +} + +func collectHistoryBatch(ctx context.Context, fetcher kafkaMessageFetcher, first kafka.Message, maxSize int, maxWait time.Duration) []kafka.Message { + if maxSize <= 1 { + return []kafka.Message{first} + } + if maxWait <= 0 { + maxWait = 100 * time.Millisecond + } + batch := []kafka.Message{first} + deadline := time.Now().Add(maxWait) + for len(batch) < maxSize { + remaining := time.Until(deadline) + if remaining <= 0 { + break + } + fetchCtx, cancel := context.WithTimeout(ctx, remaining) + message, err := fetcher.FetchMessage(fetchCtx) + cancel() + if err != nil { + break + } + batch = append(batch, message) + } + return batch +} + func processHistoryMessage(ctx context.Context, logger interface { Error(string, ...any) Warn(string, ...any) @@ -117,6 +151,64 @@ func processHistoryMessage(ctx context.Context, logger interface { addWriterMetric(registry, "vehicle_history_kafka_commits_total", message, "ok") } +func processHistoryBatch(ctx context.Context, logger interface { + Error(string, ...any) + Warn(string, ...any) +}, registry *metrics.Registry, appender historyAppender, committer kafkaMessageCommitter, messages []kafka.Message) { + if len(messages) == 0 { + return + } + messageCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), kafkaMessageOperationTimeout) + defer cancel() + + envelopes := make([]envelope.FrameEnvelope, 0, len(messages)) + for _, message := range messages { + addWriterMetric(registry, "vehicle_history_kafka_messages_total", message, "received") + addWriterLagMetric(registry, message) + var env envelope.FrameEnvelope + if err := json.Unmarshal(message.Value, &env); err != nil { + addWriterMetric(registry, "vehicle_history_kafka_messages_total", message, "invalid_json") + logger.Warn("skip invalid envelope json", "topic", message.Topic, "partition", message.Partition, "offset", message.Offset, "error", err) + continue + } + envelopes = append(envelopes, env) + } + if len(envelopes) > 0 { + started := time.Now() + err := appender.AppendAllBatch(messageCtx, envelopes) + elapsed := time.Since(started) + status := "ok" + if err != nil { + status = "error" + } + addBatchMetric(registry, "vehicle_history_batch_flush_total", status, 1) + addBatchMetric(registry, "vehicle_history_batch_rows_total", status, float64(len(envelopes))) + setBatchDuration(registry, status, elapsed) + if err != nil { + for _, message := range messages { + addWriterMetric(registry, "vehicle_history_writes_total", message, "error") + } + first := messages[0] + logger.Error("tdengine batch append failed", "topic", first.Topic, "partition", first.Partition, "offset", first.Offset, "rows", len(envelopes), "error", err) + return + } + for _, message := range messages { + addWriterMetric(registry, "vehicle_history_writes_total", message, "ok") + } + } + if err := committer.CommitMessages(messageCtx, messages...); err != nil { + for _, message := range messages { + addWriterMetric(registry, "vehicle_history_kafka_commits_total", message, "error") + } + first := messages[0] + logger.Error("kafka batch commit failed", "topic", first.Topic, "partition", first.Partition, "offset", first.Offset, "messages", len(messages), "error", err) + return + } + for _, message := range messages { + addWriterMetric(registry, "vehicle_history_kafka_commits_total", message, "ok") + } +} + func addWriterMetric(registry *metrics.Registry, name string, message kafka.Message, status string) { if registry == nil { return @@ -124,6 +216,20 @@ func addWriterMetric(registry *metrics.Registry, name string, message kafka.Mess registry.IncCounter(name, metrics.Labels{"topic": message.Topic, "status": status}) } +func addBatchMetric(registry *metrics.Registry, name string, status string, value float64) { + if registry == nil { + return + } + registry.AddCounter(name, metrics.Labels{"status": status}, value) +} + +func setBatchDuration(registry *metrics.Registry, status string, elapsed time.Duration) { + if registry == nil { + return + } + registry.SetGauge("vehicle_history_batch_flush_duration_ms", metrics.Labels{"status": status}, float64(elapsed.Milliseconds())) +} + func addWriterLagMetric(registry *metrics.Registry, message kafka.Message) { if registry == nil { return @@ -139,6 +245,8 @@ type config struct { TDengineDSN string TDengineDatabase string EnsureSchema bool + BatchSize int + BatchWait int } func loadConfig() config { @@ -150,6 +258,8 @@ func loadConfig() config { TDengineDSN: env("TDENGINE_DSN", ""), TDengineDatabase: env("TDENGINE_DATABASE", history.DefaultDatabase), EnsureSchema: env("TDENGINE_ENSURE_SCHEMA", "true") != "false", + BatchSize: envInt("HISTORY_BATCH_SIZE", 200), + BatchWait: envInt("HISTORY_BATCH_WAIT_MS", 100), } } @@ -161,6 +271,18 @@ func env(key string, fallback string) string { return value } +func envInt(key string, fallback int) int { + value := strings.TrimSpace(os.Getenv(key)) + if value == "" { + return fallback + } + parsed, err := strconv.Atoi(value) + if err != nil { + return fallback + } + return parsed +} + func splitCSV(value string) []string { var out []string for _, item := range strings.Split(value, ",") { diff --git a/go/vehicle-gateway/cmd/history-writer/main_test.go b/go/vehicle-gateway/cmd/history-writer/main_test.go index f6460f20..c911df31 100644 --- a/go/vehicle-gateway/cmd/history-writer/main_test.go +++ b/go/vehicle-gateway/cmd/history-writer/main_test.go @@ -3,6 +3,7 @@ package main import ( "context" "encoding/json" + "errors" "strings" "testing" @@ -74,6 +75,89 @@ func TestProcessHistoryMessageRecordsMetrics(t *testing.T) { } } +func TestProcessHistoryBatchAppendsAllBeforeCommit(t *testing.T) { + first := envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "VIN001", MessageID: "0x02"} + second := envelope.FrameEnvelope{Protocol: envelope.ProtocolJT808, Phone: "13307795425", MessageID: "0x0200"} + firstPayload, err := json.Marshal(first) + if err != nil { + t.Fatal(err) + } + secondPayload, err := json.Marshal(second) + if err != nil { + t.Fatal(err) + } + appender := &contextCheckingHistoryAppender{} + committer := &contextCheckingHistoryCommitter{} + registry := metrics.NewRegistry() + + processHistoryBatch( + context.Background(), + discardHistoryLogger{}, + registry, + appender, + committer, + []kafka.Message{ + {Topic: "vehicle.raw.go.gb32960.v1", Partition: 1, Offset: 10, HighWaterMark: 13, Value: firstPayload}, + {Topic: "vehicle.raw.go.jt808.v1", Partition: 2, Offset: 20, HighWaterMark: 21, Value: secondPayload}, + }, + ) + + if appender.batchCount != 1 { + t.Fatalf("batch appends = %d, want 1", appender.batchCount) + } + if got := len(appender.batch); got != 2 { + t.Fatalf("batch size = %d, want 2", got) + } + if committer.count != 1 { + t.Fatalf("commit calls = %d, want 1", committer.count) + } + if committer.messageCount != 2 { + t.Fatalf("committed messages = %d, want 2", committer.messageCount) + } + text := registry.Render() + for _, want := range []string{ + `vehicle_history_batch_rows_total{status="ok"} 2`, + `vehicle_history_batch_flush_total{status="ok"} 1`, + } { + if !strings.Contains(text, want) { + t.Fatalf("batch metric missing %s:\n%s", want, text) + } + } +} + +func TestProcessHistoryBatchDoesNotCommitWhenAppendFails(t *testing.T) { + env := envelope.FrameEnvelope{Protocol: envelope.ProtocolGB32960, VIN: "VIN001", MessageID: "0x02"} + payload, err := json.Marshal(env) + if err != nil { + t.Fatal(err) + } + appender := &contextCheckingHistoryAppender{err: errTestHistoryAppend} + committer := &contextCheckingHistoryCommitter{} + registry := metrics.NewRegistry() + + processHistoryBatch( + context.Background(), + discardHistoryLogger{}, + registry, + appender, + committer, + []kafka.Message{{Topic: "vehicle.raw.go.gb32960.v1", Partition: 1, Offset: 10, HighWaterMark: 11, Value: payload}}, + ) + + if committer.count != 0 { + t.Fatalf("commit calls = %d, want 0", committer.count) + } + text := registry.Render() + for _, want := range []string{ + `vehicle_history_batch_rows_total{status="error"} 1`, + `vehicle_history_batch_flush_total{status="error"} 1`, + } { + if !strings.Contains(text, want) { + t.Fatalf("batch error metric missing %s:\n%s", want, text) + } + } +} + func TestLoadConfigDefaultsToGoRawTopics(t *testing.T) { cfg := loadConfig() @@ -82,27 +166,65 @@ func TestLoadConfigDefaultsToGoRawTopics(t *testing.T) { if got != want { t.Fatalf("KafkaTopics = %q, want %q", got, want) } + if cfg.BatchSize != 200 { + t.Fatalf("BatchSize = %d, want 200", cfg.BatchSize) + } + if cfg.BatchWait != 100 { + t.Fatalf("BatchWait = %d, want 100", cfg.BatchWait) + } +} + +func TestLoadConfigReadsBatchSettings(t *testing.T) { + t.Setenv("HISTORY_BATCH_SIZE", "500") + t.Setenv("HISTORY_BATCH_WAIT_MS", "250") + + cfg := loadConfig() + + if cfg.BatchSize != 500 { + t.Fatalf("BatchSize = %d, want 500", cfg.BatchSize) + } + if cfg.BatchWait != 250 { + t.Fatalf("BatchWait = %d, want 250", cfg.BatchWait) + } } type contextCheckingHistoryAppender struct { - ctxErr error - count int + ctxErr error + count int + batchCount int + batch []envelope.FrameEnvelope + err error } func (a *contextCheckingHistoryAppender) AppendAll(ctx context.Context, _ envelope.FrameEnvelope) error { a.ctxErr = ctx.Err() a.count++ - return a.ctxErr + if a.ctxErr != nil { + return a.ctxErr + } + return a.err +} + +func (a *contextCheckingHistoryAppender) AppendAllBatch(ctx context.Context, envs []envelope.FrameEnvelope) error { + a.ctxErr = ctx.Err() + a.batchCount++ + a.batch = append([]envelope.FrameEnvelope(nil), envs...) + if a.ctxErr != nil { + return a.ctxErr + } + return a.err } type contextCheckingHistoryCommitter struct { - ctxErr error - count int + ctxErr error + count int + messageCount int } -func (c *contextCheckingHistoryCommitter) CommitMessages(ctx context.Context, _ ...kafka.Message) error { +func (c *contextCheckingHistoryCommitter) CommitMessages(ctx context.Context, messages ...kafka.Message) error { c.ctxErr = ctx.Err() c.count++ + c.messageCount += len(messages) return c.ctxErr } @@ -110,3 +232,5 @@ type discardHistoryLogger struct{} func (discardHistoryLogger) Error(string, ...any) {} func (discardHistoryLogger) Warn(string, ...any) {} + +var errTestHistoryAppend = errors.New("test history append failed") diff --git a/go/vehicle-gateway/internal/history/writer.go b/go/vehicle-gateway/internal/history/writer.go index aba25e25..69bce4fb 100644 --- a/go/vehicle-gateway/internal/history/writer.go +++ b/go/vehicle-gateway/internal/history/writer.go @@ -66,6 +66,16 @@ func (w *Writer) AppendAll(ctx context.Context, env envelope.FrameEnvelope) erro return w.AppendLocation(ctx, env) } +func (w *Writer) AppendAllBatch(ctx context.Context, envelopes []envelope.FrameEnvelope) error { + if len(envelopes) == 0 { + return nil + } + if err := w.AppendRawFrameBatch(ctx, envelopes); err != nil { + return err + } + return w.AppendLocationBatch(ctx, envelopes) +} + func (w *Writer) AppendRawFrame(ctx context.Context, env envelope.FrameEnvelope) error { table := tableName("raw", env) if err := w.ensureRawChild(ctx, table, "raw_frames", env); err != nil { @@ -100,6 +110,58 @@ VALUES (%s)`, chunkTable, joinLiterals(chunkValues(env, chunk)))); err != nil { return nil } +func (w *Writer) AppendRawFrameBatch(ctx context.Context, envelopes []envelope.FrameEnvelope) error { + rowsByTable := map[string][]string{} + chunkRowsByTable := map[string][]string{} + chunkEnvByTable := map[string]envelope.FrameEnvelope{} + for _, env := range envelopes { + table := tableName("raw", env) + if err := w.ensureRawChild(ctx, table, "raw_frames", env); err != nil { + return err + } + rawHex, rawHexChunks := chunkPayload(env, "raw_hex", env.RawHex) + rawText, rawTextChunks := chunkPayload(env, "raw_text", env.RawText) + parsedFields, parsedChunks := chunkPayload(env, "parsed_fields", parsedFieldsJSONString(env)) + rowsByTable[table] = append(rowsByTable[table], "("+joinLiterals(rawValues(env, rawHex, rawText, parsedFields))+")") + chunks := append(rawHexChunks, rawTextChunks...) + chunks = append(chunks, parsedChunks...) + if len(chunks) == 0 { + continue + } + chunkTable := tableName("chunk", env) + if err := w.ensureRawChild(ctx, chunkTable, "raw_frame_payload_chunks", env); err != nil { + return err + } + chunkEnvByTable[chunkTable] = env + for _, chunk := range chunks { + chunkRowsByTable[chunkTable] = append(chunkRowsByTable[chunkTable], "("+joinLiterals(chunkValues(env, chunk))+")") + } + } + for table, rows := range rowsByTable { + if len(rows) == 0 { + continue + } + if _, err := w.exec.ExecContext(ctx, fmt.Sprintf(`INSERT INTO %s + (ts, frame_id, event_id, message_id, event_time, received_at, raw_size_bytes, + raw_hex, raw_text, parsed_json, parse_status, parse_error, source_endpoint) +VALUES %s`, table, strings.Join(rows, ","))); err != nil { + return err + } + } + for table, rows := range chunkRowsByTable { + if len(rows) == 0 { + continue + } + _ = chunkEnvByTable[table] + if _, err := w.exec.ExecContext(ctx, fmt.Sprintf(`INSERT INTO %s + (ts, event_id, frame_id, received_at, payload_kind, chunk_index, chunk_count, chunk_text) +VALUES %s`, table, strings.Join(rows, ","))); err != nil { + return err + } + } + return nil +} + func (w *Writer) AppendLocation(ctx context.Context, env envelope.FrameEnvelope) error { if strings.TrimSpace(env.VIN) == "" { return nil @@ -120,6 +182,37 @@ VALUES (%s)`, table, joinLiterals(locationValues(env, longitude, latitude)))) return err } +func (w *Writer) AppendLocationBatch(ctx context.Context, envelopes []envelope.FrameEnvelope) error { + rowsByTable := map[string][]string{} + for _, env := range envelopes { + if strings.TrimSpace(env.VIN) == "" { + continue + } + longitude, okLon := floatField(env, envelope.FieldLongitude) + latitude, okLat := floatField(env, envelope.FieldLatitude) + if !okLon || !okLat { + continue + } + table := locationTableName(env) + if err := w.ensureLocationChild(ctx, table, env); err != nil { + return err + } + rowsByTable[table] = append(rowsByTable[table], "("+joinLiterals(locationValues(env, longitude, latitude))+")") + } + for table, rows := range rowsByTable { + if len(rows) == 0 { + continue + } + if _, err := w.exec.ExecContext(ctx, fmt.Sprintf(`INSERT INTO %s + (ts, event_id, received_at, longitude, latitude, altitude_m, speed_kmh, + direction_deg, alarm_flag, status_flag, total_mileage_km) +VALUES %s`, table, strings.Join(rows, ","))); err != nil { + return err + } + } + return nil +} + func (w *Writer) ensureRawChild(ctx context.Context, table string, stable string, env envelope.FrameEnvelope) error { key := stable + "." + table w.cache.mu.Lock() diff --git a/go/vehicle-gateway/internal/history/writer_test.go b/go/vehicle-gateway/internal/history/writer_test.go index 535f72c0..150f819c 100644 --- a/go/vehicle-gateway/internal/history/writer_test.go +++ b/go/vehicle-gateway/internal/history/writer_test.go @@ -226,6 +226,53 @@ func TestWriterSkipsLocationWhenVINIsMissing(t *testing.T) { } } +func TestWriterAppendsBatchRowsByChildTable(t *testing.T) { + exec := &recordingExec{} + writer := NewWriter(exec) + first := sampleEnvelope() + second := sampleEnvelope() + second.Sequence = 2 + second.EventTimeMS += 1000 + second.ReceivedAtMS += 1000 + + if err := writer.AppendAllBatch(context.Background(), []envelope.FrameEnvelope{first, second}); err != nil { + t.Fatalf("AppendAllBatch() error = %v", err) + } + + if got := countSQL(exec.calls, "USING raw_frames"); got != 1 { + t.Fatalf("raw child create count = %d", got) + } + if got := countSQL(exec.calls, "USING vehicle_locations"); got != 1 { + t.Fatalf("location child create count = %d", got) + } + if got := countSQL(exec.calls, "INSERT INTO raw_"); got != 1 { + t.Fatalf("raw batch insert count = %d", got) + } + if got := countSQL(exec.calls, "INSERT INTO loc_"); got != 1 { + t.Fatalf("location batch insert count = %d", got) + } + rawInsert := findSQL(exec.calls, "INSERT INTO raw_") + if got := strings.Count(rawInsert, "),(") + 1; got != 2 { + t.Fatalf("raw batch row count = %d, sql=%s", got, rawInsert) + } + locationInsert := findSQL(exec.calls, "INSERT INTO loc_") + if got := strings.Count(locationInsert, "),(") + 1; got != 2 { + t.Fatalf("location batch row count = %d, sql=%s", got, locationInsert) + } +} + +func TestWriterAppendAllBatchSkipsEmptyBatch(t *testing.T) { + exec := &recordingExec{} + writer := NewWriter(exec) + + if err := writer.AppendAllBatch(context.Background(), nil); err != nil { + t.Fatalf("AppendAllBatch() error = %v", err) + } + if len(exec.calls) != 0 { + t.Fatalf("exec calls = %#v, want none", exec.calls) + } +} + func TestRawSizeBytesUsesRawTextWhenHexIsEmpty(t *testing.T) { env := envelope.FrameEnvelope{RawText: `{"code":"0F80","data":{"speed":12}}`} if got, want := rawSizeBytes(env), len([]byte(env.RawText)); got != want {