From ae618905a365bcd2a988f08bb021211b17ec938f Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 2 Jul 2026 15:20:29 +0800 Subject: [PATCH] =?UTF-8?q?perf(clickhouse):=20P0=20write=20path=20?= =?UTF-8?q?=E2=80=94=20remove=20heartbeat=20DELETE,=20batchwriter=20MinBat?= =?UTF-8?q?chSize,=20tune=20chwriter?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../apps/openflare/agent/observability.go | 8 +- internal/apps/openflare/chwriter/writer.go | 22 +++- internal/db/batchwriter/config.go | 15 ++- internal/db/batchwriter/writer.go | 4 +- internal/db/batchwriter/writer_test.go | 107 ++++++++++++++++++ 5 files changed, 143 insertions(+), 13 deletions(-) diff --git a/internal/apps/openflare/agent/observability.go b/internal/apps/openflare/agent/observability.go index d210fe7f..851bd5b0 100644 --- a/internal/apps/openflare/agent/observability.go +++ b/internal/apps/openflare/agent/observability.go @@ -24,8 +24,6 @@ const ( healthSeverityInfo = "info" healthSeverityWarning = "warning" healthSeverityCritical = "critical" - nodeAccessLogRetentionDays = 90 - nodeAccessLogRetentionWindow = nodeAccessLogRetentionDays * 24 * time.Hour accessLogPathMaxLength = 100 healthEventMessageMaxLength = 4096 ) @@ -241,11 +239,7 @@ func persistNodeAccessLogs(ctx context.Context, nodeID string, records []*model. if len(records) == 0 { return nil } - if err := model.InsertOpenFlareAccessLogsBatch(ctx, records); err != nil { - return err - } - _, err := model.DeleteOpenFlareAccessLogsByNodeBefore(ctx, nodeID, reportedAt.Add(-nodeAccessLogRetentionWindow)) - return err + return model.InsertOpenFlareAccessLogsBatch(ctx, records) } func reconcileNodeHealthEvents(tx *gorm.DB, nodeID string, events []NodeHealthEvent, reportedAt time.Time) error { diff --git a/internal/apps/openflare/chwriter/writer.go b/internal/apps/openflare/chwriter/writer.go index 8e8f0a03..5a35481f 100644 --- a/internal/apps/openflare/chwriter/writer.go +++ b/internal/apps/openflare/chwriter/writer.go @@ -21,8 +21,8 @@ import ( const ( observabilityQueueSize = 5_000 - observabilityMaxBatchSize = 200 - observabilityFlushEvery = 2 * time.Second + observabilityMaxBatchSize = 500 + observabilityFlushEvery = 5 * time.Second nodeAccessLogQueueSize = 10_000 nodeAccessLogMaxBatchSize = 1_000 @@ -41,6 +41,9 @@ var ( metricSnapshotDedup *dedupSet requestReportDedup *dedupSet + openrestyDedup *dedupSet + frpsDedup *dedupSet + frpcDedup *dedupSet ) // Init starts OpenFlare ClickHouse batch writers. Safe to call multiple times. @@ -52,6 +55,9 @@ func Init(ctx context.Context) { initOnce.Do(func() { metricSnapshotDedup = newDedupSet() requestReportDedup = newDedupSet() + openrestyDedup = newDedupSet() + frpsDedup = newDedupSet() + frpcDedup = newDedupSet() metricSnapshotWriter = mustNewObservabilityWriter("metric_snapshots", analyticsrepo.BatchInsertNodeMetricSnapshots) requestReportWriter = mustNewObservabilityWriter("request_reports", analyticsrepo.BatchInsertNodeRequestReports) @@ -130,6 +136,10 @@ func QueueOpenrestyObservation(observation analyticsmodel.NodeObsOpenresty) { if openrestyWriter == nil { return } + key := fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano()) + if !openrestyDedup.markIfNew(key) { + return + } openrestyWriter.TryEnqueue(observation) } @@ -138,6 +148,10 @@ func QueueFrpsObservation(observation analyticsmodel.NodeObsFrps) { if frpsWriter == nil { return } + key := fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano()) + if !frpsDedup.markIfNew(key) { + return + } frpsWriter.TryEnqueue(observation) } @@ -146,6 +160,10 @@ func QueueFrpcObservation(observation analyticsmodel.NodeObsFrpc) { if frpcWriter == nil { return } + key := fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano()) + if !frpcDedup.markIfNew(key) { + return + } frpcWriter.TryEnqueue(observation) } diff --git a/internal/db/batchwriter/config.go b/internal/db/batchwriter/config.go index 32972ecb..b2dcfef7 100644 --- a/internal/db/batchwriter/config.go +++ b/internal/db/batchwriter/config.go @@ -9,9 +9,10 @@ import ( ) const ( - defaultQueueSize = 10_000 - defaultMaxBatchSize = 1_000 - defaultFlushEvery = time.Second + defaultQueueSize = 10_000 + defaultMaxBatchSize = 1_000 + defaultMinBatchSize = 50 + defaultFlushEvery = time.Second ) // Config controls queue capacity and flush thresholds for a Writer instance. @@ -25,6 +26,10 @@ type Config struct { // MaxBatchSize triggers a flush when the in-memory batch reaches this count. MaxBatchSize int + // MinBatchSize is the minimum in-memory batch size for time-based flushes. + // Zero disables the threshold and preserves legacy interval flush behavior. + MinBatchSize int + // FlushInterval triggers a time-based flush even when the batch is smaller. FlushInterval time.Duration } @@ -34,6 +39,7 @@ func DefaultConfig() Config { return Config{ QueueSize: defaultQueueSize, MaxBatchSize: defaultMaxBatchSize, + MinBatchSize: defaultMinBatchSize, FlushInterval: defaultFlushEvery, } } @@ -45,6 +51,9 @@ func (c Config) validate() error { if c.MaxBatchSize <= 0 { return fmt.Errorf("batchwriter: max batch size must be positive") } + if c.MinBatchSize < 0 { + return fmt.Errorf("batchwriter: min batch size must be non-negative") + } if c.FlushInterval <= 0 { return fmt.Errorf("batchwriter: flush interval must be positive") } diff --git a/internal/db/batchwriter/writer.go b/internal/db/batchwriter/writer.go index ac8c24e4..3e2672cb 100644 --- a/internal/db/batchwriter/writer.go +++ b/internal/db/batchwriter/writer.go @@ -202,7 +202,9 @@ func (w *Writer[T]) run() { flush() } case <-ticker.C: - flush() + if len(batch) > 0 && (w.cfg.MinBatchSize == 0 || len(batch) >= w.cfg.MinBatchSize) { + flush() + } } } } diff --git a/internal/db/batchwriter/writer_test.go b/internal/db/batchwriter/writer_test.go index 0fd410df..ebd80606 100644 --- a/internal/db/batchwriter/writer_test.go +++ b/internal/db/batchwriter/writer_test.go @@ -98,6 +98,7 @@ func TestWriterFlushesOnInterval(t *testing.T) { ) cfg := DefaultConfig() cfg.MaxBatchSize = 100 + cfg.MinBatchSize = 0 cfg.FlushInterval = 20 * time.Millisecond writer, err := New[int](cfg, func(_ context.Context, items []int) error { @@ -219,6 +220,112 @@ func TestWriterStopDrainsQueuedItems(t *testing.T) { } } +func TestWriterSkipsIntervalFlushBelowMinBatchSize(t *testing.T) { + t.Parallel() + + var ( + mu sync.Mutex + batch []int + ) + cfg := DefaultConfig() + cfg.MaxBatchSize = 100 + cfg.MinBatchSize = 5 + cfg.FlushInterval = 20 * time.Millisecond + + writer, err := New[int](cfg, func(_ context.Context, items []int) error { + mu.Lock() + defer mu.Unlock() + batch = append([]int(nil), items...) + return nil + }) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + writer.Start(context.Background()) + t.Cleanup(func() { + stopCtx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + if err := writer.Stop(stopCtx); err != nil { + t.Fatalf("Stop() error = %v", err) + } + }) + + for i := range 3 { + if !writer.TryEnqueue(i + 1) { + t.Fatalf("TryEnqueue(%d) = false, want true", i+1) + } + } + + time.Sleep(100 * time.Millisecond) + + mu.Lock() + got := batch + mu.Unlock() + + if len(got) != 0 { + t.Fatalf("interval flush with below-min batch = %v, want no flush", got) + } +} + +func TestWriterFlushesOnIntervalWhenMinBatchSizeReached(t *testing.T) { + t.Parallel() + + var ( + mu sync.Mutex + batch []int + ) + cfg := DefaultConfig() + cfg.MaxBatchSize = 100 + cfg.MinBatchSize = 3 + cfg.FlushInterval = 20 * time.Millisecond + + writer, err := New[int](cfg, func(_ context.Context, items []int) error { + mu.Lock() + defer mu.Unlock() + batch = append([]int(nil), items...) + return nil + }) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + writer.Start(context.Background()) + t.Cleanup(func() { + stopCtx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + if err := writer.Stop(stopCtx); err != nil { + t.Fatalf("Stop() error = %v", err) + } + }) + + for i := range 3 { + if !writer.TryEnqueue(i + 1) { + t.Fatalf("TryEnqueue(%d) = false, want true", i+1) + } + } + + deadline := time.Now().Add(time.Second) + for { + mu.Lock() + ready := len(batch) == 3 + mu.Unlock() + if ready || time.Now().After(deadline) { + break + } + time.Sleep(5 * time.Millisecond) + } + + mu.Lock() + got := batch + mu.Unlock() + + want := []int{1, 2, 3} + if diff := cmp.Diff(want, got); diff != "" { + t.Fatalf("interval flush at min batch size mismatch (-want +got):\n%s", diff) + } +} + func TestWriterInvokesFlushErrorHandler(t *testing.T) { t.Parallel()