From 9b3555c56954928a26b257ba2a2115448d95b70b Mon Sep 17 00:00:00 2001 From: ryan Date: Fri, 10 Jul 2026 10:08:32 +0800 Subject: [PATCH] fix(clickhouse): cut idle CPU from tiny parts and oversized merge pools Observability writers flushed every few seconds with MinBatchSize unset, creating constant small parts and merge load. Enable MinBatchSize with MaxFlushWait, batch access logs more aggressively, and shrink ClickHouse background pools for 3c hosts. --- docker/clickhouse/config.d/performance.xml | 19 ++++++-- docs/changelog/index.md | 1 + internal/apps/openflare/chwriter/writer.go | 14 +++++- internal/db/batchwriter/config.go | 10 +++- internal/db/batchwriter/writer.go | 20 +++++++- internal/db/batchwriter/writer_test.go | 55 ++++++++++++++++++++++ 6 files changed, 111 insertions(+), 8 deletions(-) diff --git a/docker/clickhouse/config.d/performance.xml b/docker/clickhouse/config.d/performance.xml index 29e0f617..34f4bba0 100644 --- a/docker/clickhouse/config.d/performance.xml +++ b/docker/clickhouse/config.d/performance.xml @@ -1,6 +1,17 @@ + - 50 - 8 - 2 - \ No newline at end of file + 20 + 2 + 1 + 4 + 2 + 2 + 1 + 268435456 + 0 + diff --git a/docs/changelog/index.md b/docs/changelog/index.md index da96b340..28bd5613 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -21,6 +21,7 @@ sidebar: false ### 修复 - 修复节点/仪表盘 24 小时容量、网络、磁盘 IO 趋势在 ClickHouse 限流查询下几乎为空的问题:改为基于 `of_node_metric_snapshots` / `of_node_obs_openresty` 的小时级聚合构建趋势;主机与 OpenResty 累计计数器改为按小时 delta 统计。 +- 降低静置时 ClickHouse CPU:可观测/访问日志 batchwriter 启用 `MinBatchSize` 与 `MaxFlushWait`,减少心跳小 part 写入;Docker `performance.xml` 将 `background_pool_size` 等后台线程收紧到适配 3c 小规格。 ## [v3.1.1] - 2026-07-06 diff --git a/internal/apps/openflare/chwriter/writer.go b/internal/apps/openflare/chwriter/writer.go index 5a35481f..7e5fd60b 100644 --- a/internal/apps/openflare/chwriter/writer.go +++ b/internal/apps/openflare/chwriter/writer.go @@ -20,13 +20,19 @@ import ( ) const ( + // Observability traffic is sparse (heartbeat ~10s/node). Prefer larger batches to + // cut ClickHouse parts/merges; MaxFlushWait bounds visibility lag for single-node labs. observabilityQueueSize = 5_000 observabilityMaxBatchSize = 500 - observabilityFlushEvery = 5 * time.Second + observabilityMinBatchSize = 20 + observabilityFlushEvery = 10 * time.Second + observabilityMaxFlushWait = 30 * time.Second nodeAccessLogQueueSize = 10_000 nodeAccessLogMaxBatchSize = 1_000 - nodeAccessLogFlushEvery = time.Second + nodeAccessLogMinBatchSize = 50 + nodeAccessLogFlushEvery = 2 * time.Second + nodeAccessLogMaxFlushWait = 5 * time.Second ) var ( @@ -182,7 +188,9 @@ func mustNewObservabilityWriter[T any](name string, flush batchwriter.FlushFunc[ Name: name, QueueSize: observabilityQueueSize, MaxBatchSize: observabilityMaxBatchSize, + MinBatchSize: observabilityMinBatchSize, FlushInterval: observabilityFlushEvery, + MaxFlushWait: observabilityMaxFlushWait, } writer, err := batchwriter.New( cfg, @@ -203,7 +211,9 @@ func mustNewNodeAccessLogWriter() *batchwriter.Writer[analyticsmodel.NodeAccessL Name: "node_access_logs", QueueSize: nodeAccessLogQueueSize, MaxBatchSize: nodeAccessLogMaxBatchSize, + MinBatchSize: nodeAccessLogMinBatchSize, FlushInterval: nodeAccessLogFlushEvery, + MaxFlushWait: nodeAccessLogMaxFlushWait, } writer, err := batchwriter.New[analyticsmodel.NodeAccessLog](cfg, analyticsrepo.BatchInsertNodeAccessLogs, batchwriter.WithDropHandler[analyticsmodel.NodeAccessLog](func(item analyticsmodel.NodeAccessLog) { diff --git a/internal/db/batchwriter/config.go b/internal/db/batchwriter/config.go index b2dcfef7..1c8c9f1e 100644 --- a/internal/db/batchwriter/config.go +++ b/internal/db/batchwriter/config.go @@ -28,10 +28,15 @@ type Config struct { // MinBatchSize is the minimum in-memory batch size for time-based flushes. // Zero disables the threshold and preserves legacy interval flush behavior. + // When set, interval flushes below this size are skipped unless MaxFlushWait elapses. MinBatchSize int - // FlushInterval triggers a time-based flush even when the batch is smaller. + // FlushInterval is how often the worker checks whether a time-based flush should run. FlushInterval time.Duration + + // MaxFlushWait forces a flush of any non-empty batch once the oldest item has waited + // this long, even if MinBatchSize has not been reached. Zero disables the force path. + MaxFlushWait time.Duration } // DefaultConfig returns production-friendly defaults aligned with audit log batching. @@ -57,5 +62,8 @@ func (c Config) validate() error { if c.FlushInterval <= 0 { return fmt.Errorf("batchwriter: flush interval must be positive") } + if c.MaxFlushWait < 0 { + return fmt.Errorf("batchwriter: max flush wait must be non-negative") + } return nil } \ No newline at end of file diff --git a/internal/db/batchwriter/writer.go b/internal/db/batchwriter/writer.go index 3e2672cb..13890fb4 100644 --- a/internal/db/batchwriter/writer.go +++ b/internal/db/batchwriter/writer.go @@ -173,6 +173,7 @@ func (w *Writer[T]) run() { defer ticker.Stop() batch := make([]T, 0, w.cfg.MaxBatchSize) + var batchStartedAt time.Time flush := func() { if len(batch) == 0 { return @@ -184,6 +185,7 @@ func (w *Writer[T]) run() { } } batch = batch[:0] + batchStartedAt = time.Time{} } defer func() { @@ -197,18 +199,34 @@ func (w *Writer[T]) run() { if !ok { return } + if len(batch) == 0 { + batchStartedAt = time.Now() + } batch = append(batch, item) if len(batch) >= w.cfg.MaxBatchSize { flush() } case <-ticker.C: - if len(batch) > 0 && (w.cfg.MinBatchSize == 0 || len(batch) >= w.cfg.MinBatchSize) { + if w.shouldFlushOnInterval(len(batch), batchStartedAt, time.Now()) { flush() } } } } +func (w *Writer[T]) shouldFlushOnInterval(batchLen int, batchStartedAt time.Time, now time.Time) bool { + if batchLen == 0 { + return false + } + if w.cfg.MinBatchSize == 0 || batchLen >= w.cfg.MinBatchSize { + return true + } + if w.cfg.MaxFlushWait <= 0 || batchStartedAt.IsZero() { + return false + } + return !now.Before(batchStartedAt.Add(w.cfg.MaxFlushWait)) +} + func (w *Writer[T]) notifyDrop(item T) { if w.onDrop == nil { return diff --git a/internal/db/batchwriter/writer_test.go b/internal/db/batchwriter/writer_test.go index ebd80606..13c0fb1d 100644 --- a/internal/db/batchwriter/writer_test.go +++ b/internal/db/batchwriter/writer_test.go @@ -326,6 +326,61 @@ func TestWriterFlushesOnIntervalWhenMinBatchSizeReached(t *testing.T) { } } +func TestWriterForcesFlushAfterMaxFlushWait(t *testing.T) { + t.Parallel() + + var ( + mu sync.Mutex + batch []int + ) + cfg := DefaultConfig() + cfg.MaxBatchSize = 100 + cfg.MinBatchSize = 50 + cfg.FlushInterval = 20 * time.Millisecond + cfg.MaxFlushWait = 80 * 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) + } + }) + + if !writer.TryEnqueue(1) { + t.Fatal("TryEnqueue(1) = false, want true") + } + + deadline := time.Now().Add(time.Second) + for { + mu.Lock() + ready := len(batch) == 1 + mu.Unlock() + if ready || time.Now().After(deadline) { + break + } + time.Sleep(5 * time.Millisecond) + } + + mu.Lock() + got := batch + mu.Unlock() + if diff := cmp.Diff([]int{1}, got); diff != "" { + t.Fatalf("max flush wait mismatch (-want +got):\n%s", diff) + } +} + func TestWriterInvokesFlushErrorHandler(t *testing.T) { t.Parallel()