mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-02 06:56:36 +08:00
fix(clickhouse): harden R/W path P0–P3 (cleanup, durability, rollups)
Honest TTL cleanup semantics; enqueue-safe dedup with flush retry and writer metrics; model insert hooks; latest-per-node and hourly metric/openresty rollups; small-host pool/async defaults, traffic hourly TTL, and UV labeling.
This commit is contained in:
@@ -9,22 +9,34 @@ package batchwriter
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
)
|
||||
|
||||
// FlushFunc persists a batch of queued items. It is invoked from the worker goroutine.
|
||||
type FlushFunc[T any] func(ctx context.Context, items []T) error
|
||||
|
||||
// FlushErrorHandler is called when FlushFunc returns an error. The batch is discarded
|
||||
// after the handler returns; the worker continues processing.
|
||||
type FlushErrorHandler func(ctx context.Context, batchSize int, err error)
|
||||
// FlushErrorHandler is called when FlushFunc returns an error after optional retries.
|
||||
// The batch is discarded after the handler returns; the worker continues processing.
|
||||
// Handlers receive the failed items so callers can release dedup keys or re-queue.
|
||||
type FlushErrorHandler[T any] func(ctx context.Context, items []T, err error)
|
||||
|
||||
// Stats is a point-in-time snapshot of Writer queue and failure counters.
|
||||
type Stats struct {
|
||||
Name string `json:"name"`
|
||||
Depth int `json:"depth"`
|
||||
Cap int `json:"cap"`
|
||||
Drops int64 `json:"drops"`
|
||||
FlushErrors int64 `json:"flush_errors"`
|
||||
Running bool `json:"running"`
|
||||
}
|
||||
|
||||
// Writer buffers items and flushes them by size or interval.
|
||||
type Writer[T any] struct {
|
||||
cfg Config
|
||||
flush FlushFunc[T]
|
||||
|
||||
onFlushError FlushErrorHandler
|
||||
onFlushError FlushErrorHandler[T]
|
||||
onDrop func(T)
|
||||
|
||||
startOnce sync.Once
|
||||
@@ -34,13 +46,16 @@ type Writer[T any] struct {
|
||||
ch chan T
|
||||
workerCtx context.Context
|
||||
done chan struct{}
|
||||
|
||||
drops atomic.Int64
|
||||
flushErrors atomic.Int64
|
||||
}
|
||||
|
||||
// Option configures optional Writer callbacks.
|
||||
type Option[T any] func(*Writer[T])
|
||||
|
||||
// WithFlushErrorHandler registers a callback for flush failures.
|
||||
func WithFlushErrorHandler[T any](handler FlushErrorHandler) Option[T] {
|
||||
func WithFlushErrorHandler[T any](handler FlushErrorHandler[T]) Option[T] {
|
||||
return func(w *Writer[T]) {
|
||||
w.onFlushError = handler
|
||||
}
|
||||
@@ -168,6 +183,18 @@ func (w *Writer[T]) Cap() int {
|
||||
return w.cfg.QueueSize
|
||||
}
|
||||
|
||||
// Stats returns a point-in-time snapshot of queue depth and failure counters.
|
||||
func (w *Writer[T]) Stats() Stats {
|
||||
return Stats{
|
||||
Name: w.cfg.Name,
|
||||
Depth: w.Len(),
|
||||
Cap: w.Cap(),
|
||||
Drops: w.drops.Load(),
|
||||
FlushErrors: w.flushErrors.Load(),
|
||||
Running: w.Running(),
|
||||
}
|
||||
}
|
||||
|
||||
func (w *Writer[T]) run() {
|
||||
ticker := time.NewTicker(w.cfg.FlushInterval)
|
||||
defer ticker.Stop()
|
||||
@@ -180,8 +207,9 @@ func (w *Writer[T]) run() {
|
||||
}
|
||||
items := append([]T(nil), batch...)
|
||||
if err := w.flush(w.workerCtx, items); err != nil {
|
||||
w.flushErrors.Add(1)
|
||||
if w.onFlushError != nil {
|
||||
w.onFlushError(w.workerCtx, len(items), err)
|
||||
w.onFlushError(w.workerCtx, items, err)
|
||||
}
|
||||
}
|
||||
batch = batch[:0]
|
||||
@@ -228,8 +256,9 @@ func (w *Writer[T]) shouldFlushOnInterval(batchLen int, batchStartedAt time.Time
|
||||
}
|
||||
|
||||
func (w *Writer[T]) notifyDrop(item T) {
|
||||
w.drops.Add(1)
|
||||
if w.onDrop == nil {
|
||||
return
|
||||
}
|
||||
w.onDrop(item)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -36,7 +36,7 @@ func TestWriterFlushesOnMaxBatchSize(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var (
|
||||
mu sync.Mutex
|
||||
mu sync.Mutex
|
||||
batches [][]int
|
||||
)
|
||||
cfg := DefaultConfig()
|
||||
@@ -385,23 +385,24 @@ func TestWriterInvokesFlushErrorHandler(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cfg := DefaultConfig()
|
||||
cfg.Name = "test-flush-err"
|
||||
cfg.MaxBatchSize = 1
|
||||
cfg.FlushInterval = time.Hour
|
||||
|
||||
flushErr := errors.New("flush failed")
|
||||
var (
|
||||
mu sync.Mutex
|
||||
errCount int
|
||||
batchSize int
|
||||
mu sync.Mutex
|
||||
errCount int
|
||||
gotItems []int
|
||||
)
|
||||
|
||||
writer, err := New[int](cfg, func(context.Context, []int) error {
|
||||
return flushErr
|
||||
}, WithFlushErrorHandler[int](func(_ context.Context, size int, err error) {
|
||||
}, WithFlushErrorHandler[int](func(_ context.Context, items []int, err error) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
errCount++
|
||||
batchSize = size
|
||||
gotItems = append([]int(nil), items...)
|
||||
if !errors.Is(err, flushErr) {
|
||||
t.Errorf("flush error = %v, want %v", err, flushErr)
|
||||
}
|
||||
@@ -434,13 +435,61 @@ func TestWriterInvokesFlushErrorHandler(t *testing.T) {
|
||||
|
||||
mu.Lock()
|
||||
gotCount := errCount
|
||||
gotSize := batchSize
|
||||
items := gotItems
|
||||
mu.Unlock()
|
||||
|
||||
if gotCount != 1 {
|
||||
t.Fatalf("flush error handler count = %d, want 1", gotCount)
|
||||
}
|
||||
if gotSize != 1 {
|
||||
t.Fatalf("flush error handler batch size = %d, want 1", gotSize)
|
||||
if diff := cmp.Diff([]int{7}, items); diff != "" {
|
||||
t.Fatalf("flush error handler items mismatch (-want +got):\n%s", diff)
|
||||
}
|
||||
}
|
||||
|
||||
stats := writer.Stats()
|
||||
if stats.FlushErrors != 1 {
|
||||
t.Fatalf("Stats().FlushErrors = %d, want 1", stats.FlushErrors)
|
||||
}
|
||||
if stats.Name != "test-flush-err" {
|
||||
t.Fatalf("Stats().Name = %q, want test-flush-err", stats.Name)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriterStatsTracksDrops(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cfg := DefaultConfig()
|
||||
cfg.Name = "test-drops"
|
||||
cfg.QueueSize = 1
|
||||
cfg.MaxBatchSize = 10
|
||||
cfg.FlushInterval = time.Hour
|
||||
|
||||
writer, err := New[int](cfg, func(context.Context, []int) error { 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()
|
||||
_ = writer.Stop(stopCtx)
|
||||
})
|
||||
|
||||
if !writer.TryEnqueue(1) {
|
||||
t.Fatal("TryEnqueue(1) = false, want true")
|
||||
}
|
||||
if writer.TryEnqueue(2) {
|
||||
t.Fatal("TryEnqueue(2) = true, want false")
|
||||
}
|
||||
|
||||
stats := writer.Stats()
|
||||
if stats.Drops != 1 {
|
||||
t.Fatalf("Stats().Drops = %d, want 1", stats.Drops)
|
||||
}
|
||||
if stats.Cap != 1 {
|
||||
t.Fatalf("Stats().Cap = %d, want 1", stats.Cap)
|
||||
}
|
||||
if !stats.Running {
|
||||
t.Fatal("Stats().Running = false, want true")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -16,10 +16,18 @@ import (
|
||||
)
|
||||
|
||||
const (
|
||||
clickhouseMaxExecTime = 60 // ClickHouse 最大执行时间(秒)
|
||||
clickhouseReadTimeoutFactor = 2 // ReadTimeout 为 DialTimeout 的倍数
|
||||
clickhouseMaxExecTime = 60 // ClickHouse 最大执行时间(秒)
|
||||
clickhouseReadTimeoutFactor = 2 // ReadTimeout 为 DialTimeout 的倍数
|
||||
|
||||
// async_insert 仅挂在运行时 ChConn(写路径)上,不进入 migrator OpenDB:
|
||||
// 迁移/DDL 需要同步可见结果,且不应走异步 insert 缓冲。
|
||||
//
|
||||
// 为何启用:batchwriter 仍可能在短间隔内写出相对小的块;服务端 async_insert
|
||||
// 把多次 INSERT 合并成更大 part,减轻 3c6g 上 background merge 的 CPU 压力。
|
||||
// wait_for_async_insert=1:调用方在 flush 返回前等待落盘,避免进程崩溃丢批。
|
||||
// max_data_size / busy_timeout:约 10MB 或 ~2s 触发刷出,在延迟与 part 数之间折中。
|
||||
clickhouseAsyncInsertMaxDataSize = 10_000_000
|
||||
clickhouseAsyncInsertBusyTimeoutMs = 1000
|
||||
clickhouseAsyncInsertBusyTimeoutMs = 2000
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -52,6 +60,8 @@ func init() {
|
||||
log.Println("[ClickHouse] connection established successfully")
|
||||
}
|
||||
|
||||
// buildClickHouseOptions builds the runtime native client options (queries + batch inserts).
|
||||
// Migrator uses a separate clickhouse.OpenDB path without async_insert settings.
|
||||
func buildClickHouseOptions() *clickhouse.Options {
|
||||
cfg := config.Config.ClickHouse
|
||||
|
||||
|
||||
+80
@@ -0,0 +1,80 @@
|
||||
-- +goose Up
|
||||
-- Hourly capacity rollups (avg CPU/memory + counter min/max for in-hour delta approximation).
|
||||
-- Network/disk counters are cumulative; max-min within an hour approximates that hour's delta
|
||||
-- (cross-hour continuity is intentionally approximate for dashboard trends).
|
||||
CREATE TABLE IF NOT EXISTS of_node_metric_capacity_hourly
|
||||
(
|
||||
node_id String,
|
||||
hour DateTime,
|
||||
cpu_usage_sum SimpleAggregateFunction(sum, Float64),
|
||||
cpu_usage_count SimpleAggregateFunction(sum, UInt64),
|
||||
memory_usage_sum SimpleAggregateFunction(sum, Float64),
|
||||
memory_usage_count SimpleAggregateFunction(sum, UInt64),
|
||||
network_rx_min SimpleAggregateFunction(min, Int64),
|
||||
network_rx_max SimpleAggregateFunction(max, Int64),
|
||||
network_tx_min SimpleAggregateFunction(min, Int64),
|
||||
network_tx_max SimpleAggregateFunction(max, Int64),
|
||||
disk_read_min SimpleAggregateFunction(min, Int64),
|
||||
disk_read_max SimpleAggregateFunction(max, Int64),
|
||||
disk_write_min SimpleAggregateFunction(min, Int64),
|
||||
disk_write_max SimpleAggregateFunction(max, Int64)
|
||||
)
|
||||
ENGINE = AggregatingMergeTree()
|
||||
PARTITION BY toYYYYMM(hour)
|
||||
ORDER BY (node_id, hour)
|
||||
TTL hour + INTERVAL 30 DAY;
|
||||
|
||||
CREATE MATERIALIZED VIEW IF NOT EXISTS of_node_metric_capacity_hourly_mv
|
||||
TO of_node_metric_capacity_hourly
|
||||
AS
|
||||
SELECT
|
||||
node_id,
|
||||
toStartOfHour(captured_at) AS hour,
|
||||
sum(cpu_usage_percent) AS cpu_usage_sum,
|
||||
toUInt64(count()) AS cpu_usage_count,
|
||||
sum(if(memory_total_bytes > 0, (memory_used_bytes * 100.0) / memory_total_bytes, 0)) AS memory_usage_sum,
|
||||
toUInt64(countIf(memory_total_bytes > 0)) AS memory_usage_count,
|
||||
min(network_rx_bytes) AS network_rx_min,
|
||||
max(network_rx_bytes) AS network_rx_max,
|
||||
min(network_tx_bytes) AS network_tx_min,
|
||||
max(network_tx_bytes) AS network_tx_max,
|
||||
min(disk_read_bytes) AS disk_read_min,
|
||||
max(disk_read_bytes) AS disk_read_max,
|
||||
min(disk_write_bytes) AS disk_write_min,
|
||||
max(disk_write_bytes) AS disk_write_max
|
||||
FROM of_node_metric_snapshots
|
||||
GROUP BY node_id, hour;
|
||||
|
||||
-- Hourly OpenResty counter rollups (min/max per node-hour for delta approximation).
|
||||
CREATE TABLE IF NOT EXISTS of_node_openresty_hourly
|
||||
(
|
||||
node_id String,
|
||||
hour DateTime,
|
||||
openresty_rx_min SimpleAggregateFunction(min, Int64),
|
||||
openresty_rx_max SimpleAggregateFunction(max, Int64),
|
||||
openresty_tx_min SimpleAggregateFunction(min, Int64),
|
||||
openresty_tx_max SimpleAggregateFunction(max, Int64)
|
||||
)
|
||||
ENGINE = AggregatingMergeTree()
|
||||
PARTITION BY toYYYYMM(hour)
|
||||
ORDER BY (node_id, hour)
|
||||
TTL hour + INTERVAL 30 DAY;
|
||||
|
||||
CREATE MATERIALIZED VIEW IF NOT EXISTS of_node_openresty_hourly_mv
|
||||
TO of_node_openresty_hourly
|
||||
AS
|
||||
SELECT
|
||||
node_id,
|
||||
toStartOfHour(captured_at) AS hour,
|
||||
min(openresty_rx_bytes) AS openresty_rx_min,
|
||||
max(openresty_rx_bytes) AS openresty_rx_max,
|
||||
min(openresty_tx_bytes) AS openresty_tx_min,
|
||||
max(openresty_tx_bytes) AS openresty_tx_max
|
||||
FROM of_node_obs_openresty
|
||||
GROUP BY node_id, hour;
|
||||
|
||||
-- +goose Down
|
||||
DROP VIEW IF EXISTS of_node_openresty_hourly_mv;
|
||||
DROP TABLE IF EXISTS of_node_openresty_hourly;
|
||||
DROP VIEW IF EXISTS of_node_metric_capacity_hourly_mv;
|
||||
DROP TABLE IF EXISTS of_node_metric_capacity_hourly;
|
||||
@@ -0,0 +1,42 @@
|
||||
-- +goose Up
|
||||
-- Hourly traffic rollups: 30d TTL + UV aggregation semantics.
|
||||
--
|
||||
-- unique_visitor_count on of_node_request_reports is per short report window
|
||||
-- (agent local distinct count for that window only). Summing those values in the
|
||||
-- MV (and again via SummingMergeTree part merges) invents a "true UV" number that
|
||||
-- double-counts visitors across windows. Prefer max() as a peak-window estimate;
|
||||
-- still NOT distinct visitors across the hour — UI/API must not overclaim.
|
||||
|
||||
ALTER TABLE of_node_traffic_hourly
|
||||
MODIFY TTL toDateTime(hour) + INTERVAL 30 DAY;
|
||||
|
||||
DROP VIEW IF EXISTS of_node_traffic_hourly_mv;
|
||||
|
||||
CREATE MATERIALIZED VIEW of_node_traffic_hourly_mv
|
||||
TO of_node_traffic_hourly
|
||||
AS
|
||||
SELECT
|
||||
node_id,
|
||||
toStartOfHour(window_ended_at) AS hour,
|
||||
sum(request_count) AS request_count,
|
||||
sum(error_count) AS error_count,
|
||||
-- Peak per-window UV estimate for the hour; not true cross-window distinct UV.
|
||||
max(unique_visitor_count) AS unique_visitor_count
|
||||
FROM of_node_request_reports
|
||||
GROUP BY node_id, hour;
|
||||
|
||||
-- +goose Down
|
||||
-- TTL reverse is not safe without table rewrite; restore prior MV definition only.
|
||||
DROP VIEW IF EXISTS of_node_traffic_hourly_mv;
|
||||
|
||||
CREATE MATERIALIZED VIEW of_node_traffic_hourly_mv
|
||||
TO of_node_traffic_hourly
|
||||
AS
|
||||
SELECT
|
||||
node_id,
|
||||
toStartOfHour(window_ended_at) AS hour,
|
||||
sum(request_count) AS request_count,
|
||||
sum(error_count) AS error_count,
|
||||
sum(unique_visitor_count) AS unique_visitor_count
|
||||
FROM of_node_request_reports
|
||||
GROUP BY node_id, hour;
|
||||
Reference in New Issue
Block a user