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()