mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-06 15:46:37 +08:00
perf(clickhouse): P0 write path — remove heartbeat DELETE, batchwriter MinBatchSize, tune chwriter
This commit is contained in:
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user