From 75e7c3fbe5929eadf37accb786ebc25c9c03079d Mon Sep 17 00:00:00 2001 From: lingniu Date: Thu, 2 Jul 2026 13:29:49 +0800 Subject: [PATCH] fix(go): batch durable spool replay --- go/vehicle-gateway/cmd/gateway/main.go | 5 +-- .../internal/eventbus/durable_sink.go | 32 +++++++++++++++---- .../internal/eventbus/durable_sink_test.go | 20 ++++++++++++ 3 files changed, 48 insertions(+), 9 deletions(-) diff --git a/go/vehicle-gateway/cmd/gateway/main.go b/go/vehicle-gateway/cmd/gateway/main.go index 246c4f4d..883161c1 100644 --- a/go/vehicle-gateway/cmd/gateway/main.go +++ b/go/vehicle-gateway/cmd/gateway/main.go @@ -169,12 +169,13 @@ func buildSink(ctx context.Context, logger *slog.Logger) (eventbus.Sink, error) }) spoolDir := strings.TrimSpace(os.Getenv("KAFKA_SPOOL_DIR")) if spoolDir != "" { - durable := eventbus.NewDurableSink(out, eventbus.DurableConfig{Directory: spoolDir}) + replayBatchSize := envInt("KAFKA_SPOOL_REPLAY_BATCH_SIZE", 1000) + durable := eventbus.NewDurableSink(out, eventbus.DurableConfig{Directory: spoolDir, ReplayBatchSize: replayBatchSize}) 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()) + logger.Info("kafka durable spool enabled", "dir", spoolDir, "replay_interval_ms", interval.Milliseconds(), "replay_batch_size", replayBatchSize) out = durable } return out, nil diff --git a/go/vehicle-gateway/internal/eventbus/durable_sink.go b/go/vehicle-gateway/internal/eventbus/durable_sink.go index eab71a14..39d8d5b8 100644 --- a/go/vehicle-gateway/internal/eventbus/durable_sink.go +++ b/go/vehicle-gateway/internal/eventbus/durable_sink.go @@ -15,7 +15,8 @@ import ( ) type DurableConfig struct { - Directory string + Directory string + ReplayBatchSize int } type DurableSink struct { @@ -25,6 +26,7 @@ type DurableSink struct { mu sync.Mutex seq uint64 rawPending map[string]struct{} + replayBatchSize int } type durableRecord struct { @@ -44,9 +46,10 @@ func NewDurableSink(delegate Sink, cfg DurableConfig) *DurableSink { panic("durable delegate sink must not be nil") } return &DurableSink{ - delegate: delegate, - dir: strings.TrimSpace(cfg.Directory), - rawPending: map[string]struct{}{}, + delegate: delegate, + dir: strings.TrimSpace(cfg.Directory), + rawPending: map[string]struct{}{}, + replayBatchSize: cfg.ReplayBatchSize, } } @@ -72,11 +75,14 @@ func (s *DurableSink) PublishUnified(ctx context.Context, env envelope.FrameEnve } func (s *DurableSink) ReplayOnce(ctx context.Context) error { - files, err := filepath.Glob(filepath.Join(s.dir, "*.json")) + return s.replay(ctx, 0) +} + +func (s *DurableSink) replay(ctx context.Context, limit int) error { + files, err := durableFiles(s.dir, limit) if err != nil { return err } - sort.Strings(files) records := make([]durableRecordFile, 0, len(files)) for _, file := range files { record, err := readDurableRecord(file) @@ -142,7 +148,7 @@ func (s *DurableSink) ReplayLoop(ctx context.Context, interval time.Duration, on func (s *DurableSink) replayOnceFromLoop(ctx context.Context) error { replayCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), durableReplayOperationTimeout) defer cancel() - return s.ReplayOnce(replayCtx) + return s.replay(replayCtx, s.replayBatchSize) } func (s *DurableSink) Close() error { @@ -224,3 +230,15 @@ func readDurableRecord(path string) (durableRecord, error) { } return record, json.Unmarshal(payload, &record) } + +func durableFiles(dir string, limit int) ([]string, error) { + files, err := filepath.Glob(filepath.Join(dir, "*.json")) + if err != nil { + return nil, err + } + sort.Strings(files) + if limit > 0 && len(files) > limit { + files = files[:limit] + } + return files, nil +} diff --git a/go/vehicle-gateway/internal/eventbus/durable_sink_test.go b/go/vehicle-gateway/internal/eventbus/durable_sink_test.go index cdf508c3..866fd37d 100644 --- a/go/vehicle-gateway/internal/eventbus/durable_sink_test.go +++ b/go/vehicle-gateway/internal/eventbus/durable_sink_test.go @@ -118,6 +118,26 @@ func TestDurableSinkReplayLoopTickUsesUncancelledContext(t *testing.T) { } } +func TestDurableSinkReplayLoopTickRespectsBatchSize(t *testing.T) { + dir := t.TempDir() + env := durableTestEnvelope() + for index := 0; index < 5; index++ { + writeDurableRecord(t, filepath.Join(dir, "000"+string(rune('1'+index))+"-raw.json"), durableRecord{Kind: "raw", Envelope: env}) + } + delegate := &contextCheckingReplaySink{} + sink := NewDurableSink(delegate, DurableConfig{Directory: dir, ReplayBatchSize: 2}) + + if err := sink.replayOnceFromLoop(context.Background()); err != nil { + t.Fatalf("replayOnceFromLoop() error = %v", err) + } + if delegate.rawCalls != 2 { + t.Fatalf("raw calls = %d, want 2", delegate.rawCalls) + } + if files := spoolFiles(t, dir); len(files) != 3 { + t.Fatalf("spool files after replay = %d, want 3: %#v", len(files), files) + } +} + func durableTestEnvelope() envelope.FrameEnvelope { return envelope.FrameEnvelope{ Protocol: envelope.ProtocolJT808,