From 243de5c2b2aa4f75954fa41fba99d9c6560cef1f Mon Sep 17 00:00:00 2001 From: lingniu Date: Wed, 1 Jul 2026 23:00:06 +0800 Subject: [PATCH] feat: spool gateway publishes to disk --- .../2026-07-01-go-ingest-redesign-design.md | 6 +- go/vehicle-gateway/cmd/gateway/main.go | 19 +- .../internal/eventbus/durable_sink.go | 187 ++++++++++++++++++ .../internal/eventbus/durable_sink_test.go | 153 ++++++++++++++ 4 files changed, 360 insertions(+), 5 deletions(-) create mode 100644 go/vehicle-gateway/internal/eventbus/durable_sink.go create mode 100644 go/vehicle-gateway/internal/eventbus/durable_sink_test.go diff --git a/docs/superpowers/specs/2026-07-01-go-ingest-redesign-design.md b/docs/superpowers/specs/2026-07-01-go-ingest-redesign-design.md index dc46d2a2..bf84986d 100644 --- a/docs/superpowers/specs/2026-07-01-go-ingest-redesign-design.md +++ b/docs/superpowers/specs/2026-07-01-go-ingest-redesign-design.md @@ -186,7 +186,9 @@ go/vehicle-gateway/ - Kafka 生产端必须开启同步写入,RAW 写成功后才允许写 unified event。 - Kafka 写入必须支持可配置重试、单次写超时和短退避,默认值为 `KAFKA_PUBLISH_ATTEMPTS=3`、`KAFKA_PUBLISH_TIMEOUT_MS=3000`、`KAFKA_PUBLISH_BACKOFF_MS=100`。 -- 当前阶段的重试只解决短暂网络抖动;后续阶段需要增加本地磁盘 spool/WAL,使 Kafka 长时间不可用时 RAW 不丢失,并支持恢复后补发。 +- gateway 支持本地磁盘 spool/WAL,配置 `KAFKA_SPOOL_DIR` 后启用;Kafka 长时间不可用时,发布失败的 envelope 先原子写入本地 JSON 文件。 +- spool 补发按文件名顺序执行,补发成功后删除文件;默认补发间隔由 `KAFKA_SPOOL_REPLAY_INTERVAL_MS=1000` 控制。 +- 如果 RAW 写 Kafka 失败并已落盘,同一个 event 的 unified event 不允许抢先写 Kafka,也必须进入 spool,确保恢复后按 RAW -> unified 顺序补发。 - Kafka topic 不允许自动创建,topic 和分区数由部署脚本或运维初始化,避免生产拼写错误造成隐性分流。 ## TDengine 数据库设计 @@ -407,6 +409,7 @@ flowchart LR - 协议坏帧仍写 RAW topic,`parse_status=BAD_FRAME`。 - 可部分解析的帧写 `parse_status=PARTIAL`,保留 parse error。 - Kafka 写失败时接入层先按配置重试;重试耗尽后不写 unified event,必须打错误日志和 metrics。 +- 启用 `KAFKA_SPOOL_DIR` 后,Kafka 重试耗尽会落本地 spool;如果 spool 写失败,才视为接入层最终失败。 - TDengine 写失败不提交 Kafka offset。 - MySQL 统计写失败不提交 Kafka offset。 - Redis 写失败不影响历史和统计,但要通过 metrics 暴露。 @@ -473,6 +476,7 @@ TDengine 连接: - Go gateway 能在 ECS 上接收真实 32960、808、宇通 MQTT 数据。 - 三种协议 RAW 都进入 Kafka。 - Kafka 短暂写失败时,gateway 按配置重试;单元测试覆盖首次失败后成功和重试耗尽返回错误。 +- Kafka 长时间不可用时,gateway 可把 RAW/unified 写入本地 spool;单元测试覆盖 RAW 已落盘时 unified 也落盘、恢复后按顺序补发并删除文件。 - 三种协议 RAW 和 parsed JSON 都进入 TDengine。 - 32960 和 808 的位置进入 `vehicle_locations`。 - 32960 和 808 的总里程采样进入 `vehicle_mileage_points`。 diff --git a/go/vehicle-gateway/cmd/gateway/main.go b/go/vehicle-gateway/cmd/gateway/main.go index 26c9991d..9e70e021 100644 --- a/go/vehicle-gateway/cmd/gateway/main.go +++ b/go/vehicle-gateway/cmd/gateway/main.go @@ -27,7 +27,7 @@ func main() { ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() - sink, err := buildSink(logger) + sink, err := buildSink(ctx, logger) if err != nil { logger.Error("build sink failed", "error", err) os.Exit(1) @@ -132,7 +132,7 @@ func buildIdentityResolver(ctx context.Context, logger *slog.Logger) (identity.R return identity.NewMySQLResolver(db, table), func() { _ = db.Close() }, nil } -func buildSink(logger *slog.Logger) (eventbus.Sink, error) { +func buildSink(ctx context.Context, logger *slog.Logger) (eventbus.Sink, error) { brokers := splitCSV(os.Getenv("KAFKA_BROKERS")) if len(brokers) == 0 { logger.Warn("KAFKA_BROKERS is empty; using log sink") @@ -150,11 +150,22 @@ func buildSink(logger *slog.Logger) (eventbus.Sink, error) { if err != nil { return nil, err } - return eventbus.NewRetryingSink(sink, eventbus.RetryConfig{ + var out eventbus.Sink = eventbus.NewRetryingSink(sink, eventbus.RetryConfig{ Attempts: envInt("KAFKA_PUBLISH_ATTEMPTS", 3), Backoff: time.Duration(envInt("KAFKA_PUBLISH_BACKOFF_MS", 100)) * time.Millisecond, AttemptTimeout: time.Duration(envInt("KAFKA_PUBLISH_TIMEOUT_MS", 3000)) * time.Millisecond, - }), nil + }) + spoolDir := strings.TrimSpace(os.Getenv("KAFKA_SPOOL_DIR")) + if spoolDir != "" { + durable := eventbus.NewDurableSink(out, eventbus.DurableConfig{Directory: spoolDir}) + interval := time.Duration(envInt("KAFKA_SPOOL_REPLAY_INTERVAL_MS", 1000)) * time.Millisecond + go durable.ReplayLoop(ctx, interval, func(err error) { + logger.Warn("kafka spool replay failed", "error", err) + }) + logger.Info("kafka durable spool enabled", "dir", spoolDir, "replay_interval_ms", interval.Milliseconds()) + out = durable + } + return out, nil } func env(key string, fallback string) string { diff --git a/go/vehicle-gateway/internal/eventbus/durable_sink.go b/go/vehicle-gateway/internal/eventbus/durable_sink.go new file mode 100644 index 00000000..c104d60d --- /dev/null +++ b/go/vehicle-gateway/internal/eventbus/durable_sink.go @@ -0,0 +1,187 @@ +package eventbus + +import ( + "context" + "encoding/json" + "fmt" + "os" + "path/filepath" + "sort" + "strings" + "sync" + "time" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +type DurableConfig struct { + Directory string +} + +type DurableSink struct { + delegate Sink + dir string + + mu sync.Mutex + seq uint64 + rawPending map[string]struct{} +} + +type durableRecord struct { + Kind string `json:"kind"` + Envelope envelope.FrameEnvelope `json:"envelope"` +} + +func NewDurableSink(delegate Sink, cfg DurableConfig) *DurableSink { + if delegate == nil { + panic("durable delegate sink must not be nil") + } + return &DurableSink{ + delegate: delegate, + dir: strings.TrimSpace(cfg.Directory), + rawPending: map[string]struct{}{}, + } +} + +func (s *DurableSink) PublishRaw(ctx context.Context, env envelope.FrameEnvelope) error { + if err := s.delegate.PublishRaw(ctx, env); err == nil { + return nil + } + if err := s.spool("raw", env); err != nil { + return err + } + s.markRawPending(env) + return nil +} + +func (s *DurableSink) PublishUnified(ctx context.Context, env envelope.FrameEnvelope) error { + if s.isRawPending(env) { + return s.spool("unified", env) + } + if err := s.delegate.PublishUnified(ctx, env); err == nil { + return nil + } + return s.spool("unified", env) +} + +func (s *DurableSink) ReplayOnce(ctx context.Context) error { + files, err := filepath.Glob(filepath.Join(s.dir, "*.json")) + if err != nil { + return err + } + sort.Strings(files) + for _, file := range files { + record, err := readDurableRecord(file) + if err != nil { + return err + } + if err := s.publishRecord(ctx, record); err != nil { + return err + } + if err := os.Remove(file); err != nil { + return err + } + if record.Kind == "raw" { + s.clearRawPending(record.Envelope) + } + } + return nil +} + +func (s *DurableSink) ReplayLoop(ctx context.Context, interval time.Duration, onError func(error)) { + if interval <= 0 { + interval = time.Second + } + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + if err := s.ReplayOnce(ctx); err != nil && onError != nil { + onError(err) + } + } + } +} + +func (s *DurableSink) Close() error { + return s.delegate.Close() +} + +func (s *DurableSink) publishRecord(ctx context.Context, record durableRecord) error { + switch record.Kind { + case "raw": + return s.delegate.PublishRaw(ctx, record.Envelope) + case "unified": + return s.delegate.PublishUnified(ctx, record.Envelope) + default: + return fmt.Errorf("unknown durable record kind %q", record.Kind) + } +} + +func (s *DurableSink) spool(kind string, env envelope.FrameEnvelope) error { + if s.dir == "" { + return fmt.Errorf("durable spool directory is empty") + } + if env.EventID == "" { + env.EventID = env.StableEventID() + } + if env.ParseStatus == "" { + env.ParseStatus = envelope.ParseOK + } + if err := os.MkdirAll(s.dir, 0o750); err != nil { + return err + } + payload, err := json.Marshal(durableRecord{Kind: kind, Envelope: env}) + if err != nil { + return err + } + name := s.nextFileName(env, kind) + path := filepath.Join(s.dir, name) + tmp := path + ".tmp" + if err := os.WriteFile(tmp, payload, 0o640); err != nil { + return err + } + return os.Rename(tmp, path) +} + +func (s *DurableSink) nextFileName(env envelope.FrameEnvelope, kind string) string { + s.mu.Lock() + defer s.mu.Unlock() + s.seq++ + eventID := env.StableEventID() + if len(eventID) > 12 { + eventID = eventID[:12] + } + return fmt.Sprintf("%020d-%06d-%s-%s.json", time.Now().UnixNano(), s.seq, kind, eventID) +} + +func (s *DurableSink) markRawPending(env envelope.FrameEnvelope) { + s.mu.Lock() + defer s.mu.Unlock() + s.rawPending[env.StableEventID()] = struct{}{} +} + +func (s *DurableSink) clearRawPending(env envelope.FrameEnvelope) { + s.mu.Lock() + defer s.mu.Unlock() + delete(s.rawPending, env.StableEventID()) +} + +func (s *DurableSink) isRawPending(env envelope.FrameEnvelope) bool { + s.mu.Lock() + defer s.mu.Unlock() + _, ok := s.rawPending[env.StableEventID()] + return ok +} + +func readDurableRecord(path string) (durableRecord, error) { + var record durableRecord + payload, err := os.ReadFile(path) + if err != nil { + return record, err + } + return record, json.Unmarshal(payload, &record) +} diff --git a/go/vehicle-gateway/internal/eventbus/durable_sink_test.go b/go/vehicle-gateway/internal/eventbus/durable_sink_test.go new file mode 100644 index 00000000..0c0f87d8 --- /dev/null +++ b/go/vehicle-gateway/internal/eventbus/durable_sink_test.go @@ -0,0 +1,153 @@ +package eventbus + +import ( + "context" + "errors" + "os" + "path/filepath" + "testing" + + "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" +) + +func TestDurableSinkSpoolsUnifiedWhenRawWasSpooled(t *testing.T) { + dir := t.TempDir() + delegate := &scriptedSink{rawErrors: []error{errSpoolTest}} + sink := NewDurableSink(delegate, DurableConfig{Directory: dir}) + env := durableTestEnvelope() + + if err := sink.PublishRaw(context.Background(), env); err != nil { + t.Fatalf("PublishRaw() error = %v", err) + } + if err := sink.PublishUnified(context.Background(), env); err != nil { + t.Fatalf("PublishUnified() error = %v", err) + } + + if delegate.rawCalls != 1 { + t.Fatalf("raw calls = %d, want 1", delegate.rawCalls) + } + if delegate.unifiedCalls != 0 { + t.Fatalf("unified should not be delegated while raw is spooled, calls = %d", delegate.unifiedCalls) + } + files := spoolFiles(t, dir) + if len(files) != 2 { + t.Fatalf("spool files = %d, want 2: %#v", len(files), files) + } + assertSpoolKind(t, files[0], "raw") + assertSpoolKind(t, files[1], "unified") +} + +func TestDurableSinkReplayPublishesInFileOrderAndDeletesFiles(t *testing.T) { + dir := t.TempDir() + delegate := &scriptedSink{rawErrors: []error{errSpoolTest}} + sink := NewDurableSink(delegate, DurableConfig{Directory: dir}) + env := durableTestEnvelope() + + if err := sink.PublishRaw(context.Background(), env); err != nil { + t.Fatalf("PublishRaw() error = %v", err) + } + if err := sink.PublishUnified(context.Background(), env); err != nil { + t.Fatalf("PublishUnified() error = %v", err) + } + delegate.rawErrors = nil + + if err := sink.ReplayOnce(context.Background()); err != nil { + t.Fatalf("ReplayOnce() error = %v", err) + } + + got := delegate.calls + want := []string{"raw", "raw", "unified"} + if len(got) != len(want) { + t.Fatalf("calls = %#v, want %#v", got, want) + } + for i := range want { + if got[i] != want[i] { + t.Fatalf("calls[%d] = %q, want %q; all=%#v", i, got[i], want[i], got) + } + } + if files := spoolFiles(t, dir); len(files) != 0 { + t.Fatalf("spool files after replay = %#v, want none", files) + } +} + +func durableTestEnvelope() envelope.FrameEnvelope { + return envelope.FrameEnvelope{ + Protocol: envelope.ProtocolJT808, + MessageID: "0x0200", + Phone: "013079963379", + VIN: "LKLG7C4E3NA774736", + Sequence: 183, + EventTimeMS: 1782903940000, + ReceivedAtMS: 1782917751903, + RawHex: "0200002201307996337900B7", + } +} + +func spoolFiles(t *testing.T, dir string) []string { + t.Helper() + matches, err := filepath.Glob(filepath.Join(dir, "*.json")) + if err != nil { + t.Fatalf("glob spool files: %v", err) + } + return matches +} + +func assertSpoolKind(t *testing.T, path string, kind string) { + t.Helper() + payload, err := os.ReadFile(path) + if err != nil { + t.Fatalf("read spool file: %v", err) + } + if !containsString(string(payload), `"kind":"`+kind+`"`) { + t.Fatalf("spool file %s payload = %s, want kind %q", path, payload, kind) + } +} + +func containsString(value string, needle string) bool { + return len(needle) == 0 || (len(value) >= len(needle) && indexString(value, needle) >= 0) +} + +func indexString(value string, needle string) int { + for i := 0; i+len(needle) <= len(value); i++ { + if value[i:i+len(needle)] == needle { + return i + } + } + return -1 +} + +var errSpoolTest = errors.New("delegate unavailable") + +type scriptedSink struct { + rawErrors []error + unifiedErrors []error + rawCalls int + unifiedCalls int + calls []string +} + +func (s *scriptedSink) PublishRaw(context.Context, envelope.FrameEnvelope) error { + s.rawCalls++ + s.calls = append(s.calls, "raw") + if len(s.rawErrors) > 0 { + err := s.rawErrors[0] + s.rawErrors = s.rawErrors[1:] + return err + } + return nil +} + +func (s *scriptedSink) PublishUnified(context.Context, envelope.FrameEnvelope) error { + s.unifiedCalls++ + s.calls = append(s.calls, "unified") + if len(s.unifiedErrors) > 0 { + err := s.unifiedErrors[0] + s.unifiedErrors = s.unifiedErrors[1:] + return err + } + return nil +} + +func (s *scriptedSink) Close() error { + return nil +}