From da1dd9240447c16cb50154a76e35b575d95d6ed8 Mon Sep 17 00:00:00 2001 From: ryan Date: Fri, 10 Jul 2026 11:24:31 +0800 Subject: [PATCH] fix(observability): merge rollup+raw hourly trends and backfill history Prefer capacity/openresty rollups only when they cover the 24h window; otherwise merge per hour so raw fills pre-MV gaps and rollup wins on overlap. Add a one-time ANTI JOIN backfill migration for the last 30 days. --- docs/changelog/index.md | 2 +- ...00003_backfill_metric_openresty_hourly.sql | 66 +++++++++++ .../analytics/node_observability.go | 112 +++++++++++++++--- .../node_observability_latest_test.go | 33 +++++- 4 files changed, 191 insertions(+), 22 deletions(-) create mode 100644 internal/db/migrator/goose/clickhouse/202607100003_backfill_metric_openresty_hourly.sql diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 1c8a8551..255ea6a3 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -20,7 +20,7 @@ sidebar: false ### 修复 -- 修复小时预聚合表仅有迁移后少量数据时仍优先于全量 raw 聚合、导致 24 小时容量/网络/磁盘趋势只显示最近时段的问题:rollup 未覆盖查询窗口起点时回退到原始快照小时聚合。 +- 修复小时预聚合仅有迁移后少量数据时 24 小时容量/网络/磁盘趋势残缺的问题:读路径按小时 merge(rollup 覆盖不足时用 raw 补洞;窗口完整时仅走 rollup);并增加历史 backfill 迁移避免长期依赖 raw。 ## [v3.1.2] - 2026-07-10 diff --git a/internal/db/migrator/goose/clickhouse/202607100003_backfill_metric_openresty_hourly.sql b/internal/db/migrator/goose/clickhouse/202607100003_backfill_metric_openresty_hourly.sql new file mode 100644 index 00000000..eb354c48 --- /dev/null +++ b/internal/db/migrator/goose/clickhouse/202607100003_backfill_metric_openresty_hourly.sql @@ -0,0 +1,66 @@ +-- +goose Up +-- One-time historical backfill for hours not yet present in rollup tables. +-- MV only ingests rows after creation; without this, 24h charts rely on raw merge forever. +-- ANTI JOIN avoids double-counting hours already filled by the live MV. + +INSERT INTO of_node_metric_capacity_hourly +SELECT + s.node_id, + toStartOfHour(s.captured_at) AS hour, + sum(s.cpu_usage_percent) AS cpu_usage_sum, + toUInt64(count()) AS cpu_usage_count, + sum(if(s.memory_total_bytes > 0, (s.memory_used_bytes * 100.0) / s.memory_total_bytes, 0)) AS memory_usage_sum, + toUInt64(countIf(s.memory_total_bytes > 0)) AS memory_usage_count, + min(s.network_rx_bytes) AS network_rx_min, + max(s.network_rx_bytes) AS network_rx_max, + min(s.network_tx_bytes) AS network_tx_min, + max(s.network_tx_bytes) AS network_tx_max, + min(s.disk_read_bytes) AS disk_read_min, + max(s.disk_read_bytes) AS disk_read_max, + min(s.disk_write_bytes) AS disk_write_min, + max(s.disk_write_bytes) AS disk_write_max +FROM of_node_metric_snapshots AS s +ANTI JOIN +( + SELECT + node_id, + hour + FROM of_node_metric_capacity_hourly + GROUP BY + node_id, + hour +) AS existing +ON s.node_id = existing.node_id AND toStartOfHour(s.captured_at) = existing.hour +WHERE s.captured_at >= now() - INTERVAL 30 DAY +GROUP BY + s.node_id, + hour; + +INSERT INTO of_node_openresty_hourly +SELECT + s.node_id, + toStartOfHour(s.captured_at) AS hour, + min(s.openresty_rx_bytes) AS openresty_rx_min, + max(s.openresty_rx_bytes) AS openresty_rx_max, + min(s.openresty_tx_bytes) AS openresty_tx_min, + max(s.openresty_tx_bytes) AS openresty_tx_max +FROM of_node_obs_openresty AS s +ANTI JOIN +( + SELECT + node_id, + hour + FROM of_node_openresty_hourly + GROUP BY + node_id, + hour +) AS existing +ON s.node_id = existing.node_id AND toStartOfHour(s.captured_at) = existing.hour +WHERE s.captured_at >= now() - INTERVAL 30 DAY +GROUP BY + s.node_id, + hour; + +-- +goose Down +-- Backfill is additive; down does not remove historical rollup rows (TTL still applies). +SELECT 1; diff --git a/internal/repository/analytics/node_observability.go b/internal/repository/analytics/node_observability.go index 467ad375..4a92c209 100644 --- a/internal/repository/analytics/node_observability.go +++ b/internal/repository/analytics/node_observability.go @@ -6,6 +6,7 @@ package analytics import ( "context" "fmt" + "sort" "time" "github.com/ClickHouse/clickhouse-go/v2/lib/driver" @@ -369,9 +370,7 @@ ORDER BY hour ASC`, nodeTrafficHourlyTableName, clause) } // hourlyRollupMaxLead is how far after filter.Since the earliest rollup bucket may start -// while still preferring the pre-aggregated tables. Materialized views only receive rows -// written after the MV exists; partial rollups (e.g. only the last hour) must not hide -// full-window raw aggregation for historical 24h charts. +// while still treating pre-aggregated tables as a complete window (skip raw query). const hourlyRollupMaxLead = 2 * time.Hour // hourlyRollupCoversWindow reports whether rollup coverage starts near the requested window. @@ -386,15 +385,61 @@ func hourlyRollupCoversWindow(earliestHour time.Time, since time.Time) bool { } // ListNodeMetricHourly returns hourly metric snapshot aggregates matching filter. -// Prefers of_node_metric_capacity_hourly when it covers the requested window; otherwise -// falls back to raw lagInFrame over of_node_metric_snapshots. +// +// Strategy (optimal for correctness + cost): +// 1. Load of_node_metric_capacity_hourly rollup. +// 2. If rollup spans the window from filter.Since, return it alone (cheap path). +// 3. Otherwise load raw lagInFrame aggregates and merge by hour: rollup wins on +// overlap, raw fills historical gaps (MV never backfills pre-creation data). func ListNodeMetricHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeMetricHourly, error) { - if rows, err := listNodeMetricHourlyFromRollup(ctx, filter); err == nil && len(rows) > 0 { - if hourlyRollupCoversWindow(rows[0].Hour, filter.Since) { - return rows, nil + rollup, rollupErr := listNodeMetricHourlyFromRollup(ctx, filter) + if rollupErr == nil && len(rollup) > 0 && hourlyRollupCoversWindow(rollup[0].Hour, filter.Since) { + return rollup, nil + } + + raw, rawErr := listNodeMetricHourlyFromRaw(ctx, filter) + if rawErr != nil { + if rollupErr == nil && len(rollup) > 0 { + return rollup, nil + } + return nil, rawErr + } + if len(rollup) == 0 { + return raw, nil + } + // Partial rollup (or rollupErr with empty slice): merge; raw fills historical gaps. + return mergeNodeMetricHourlyPreferRollup(rollup, raw), nil +} + +// mergeNodeMetricHourlyPreferRollup unions two hour series (both ASC by Hour). +// Rollup values replace raw for the same hour; raw supplies missing hours. +func mergeNodeMetricHourlyPreferRollup(rollup, raw []NodeMetricHourly) []NodeMetricHourly { + byHour := make(map[int64]NodeMetricHourly, len(raw)+len(rollup)) + order := make([]int64, 0, len(raw)+len(rollup)) + add := func(row NodeMetricHourly, overwrite bool) { + key := row.Hour.UTC().Truncate(time.Hour).Unix() + if _, exists := byHour[key]; !exists { + order = append(order, key) + byHour[key] = row + return + } + if overwrite { + byHour[key] = row } } - return listNodeMetricHourlyFromRaw(ctx, filter) + for _, row := range raw { + add(row, false) + } + for _, row := range rollup { + add(row, true) + } + result := make([]NodeMetricHourly, 0, len(order)) + // Keep chronological order of first-seen keys; re-sort by hour for stability. + sort.Slice(order, func(i, j int) bool { return order[i] < order[j] }) + for _, key := range order { + result = append(result, byHour[key]) + } + return result } func listNodeMetricHourlyFromRollup(ctx context.Context, filter NodeObservabilityFilter) ([]NodeMetricHourly, error) { @@ -512,15 +557,52 @@ func scanNodeMetricHourlyRows(rows driver.Rows) ([]NodeMetricHourly, error) { } // ListNodeOpenrestyHourly returns hourly OpenResty observation aggregates matching filter. -// Prefers of_node_openresty_hourly when it covers the requested window; otherwise -// falls back to raw lagInFrame over of_node_obs_openresty. +// Same rollup-first / per-hour merge strategy as ListNodeMetricHourly. func ListNodeOpenrestyHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) { - if rows, err := listNodeOpenrestyHourlyFromRollup(ctx, filter); err == nil && len(rows) > 0 { - if hourlyRollupCoversWindow(rows[0].Hour, filter.Since) { - return rows, nil + rollup, rollupErr := listNodeOpenrestyHourlyFromRollup(ctx, filter) + if rollupErr == nil && len(rollup) > 0 && hourlyRollupCoversWindow(rollup[0].Hour, filter.Since) { + return rollup, nil + } + + raw, rawErr := listNodeOpenrestyHourlyFromRaw(ctx, filter) + if rawErr != nil { + if rollupErr == nil && len(rollup) > 0 { + return rollup, nil + } + return nil, rawErr + } + if len(rollup) == 0 { + return raw, nil + } + return mergeNodeOpenrestyHourlyPreferRollup(rollup, raw), nil +} + +func mergeNodeOpenrestyHourlyPreferRollup(rollup, raw []NodeOpenrestyHourly) []NodeOpenrestyHourly { + byHour := make(map[int64]NodeOpenrestyHourly, len(raw)+len(rollup)) + order := make([]int64, 0, len(raw)+len(rollup)) + add := func(row NodeOpenrestyHourly, overwrite bool) { + key := row.Hour.UTC().Truncate(time.Hour).Unix() + if _, exists := byHour[key]; !exists { + order = append(order, key) + byHour[key] = row + return + } + if overwrite { + byHour[key] = row } } - return listNodeOpenrestyHourlyFromRaw(ctx, filter) + for _, row := range raw { + add(row, false) + } + for _, row := range rollup { + add(row, true) + } + sort.Slice(order, func(i, j int) bool { return order[i] < order[j] }) + result := make([]NodeOpenrestyHourly, 0, len(order)) + for _, key := range order { + result = append(result, byHour[key]) + } + return result } func listNodeOpenrestyHourlyFromRollup(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) { diff --git a/internal/repository/analytics/node_observability_latest_test.go b/internal/repository/analytics/node_observability_latest_test.go index 03c89dc0..7178f009 100644 --- a/internal/repository/analytics/node_observability_latest_test.go +++ b/internal/repository/analytics/node_observability_latest_test.go @@ -80,7 +80,7 @@ func TestListNodeMetricHourly_PrefersRollup(t *testing.T) { assert.Contains(t, mock.queries[0], nodeMetricCapacityHourlyTableName()) } -func TestListNodeMetricHourly_FallsBackWhenRollupOnlyCoversRecentHours(t *testing.T) { +func TestListNodeMetricHourly_MergesRawGapsWithPartialRollup(t *testing.T) { ctx := context.Background() // 24h window starts far before the only rollup bucket (last hour). since := time.Date(2026, 7, 9, 12, 0, 0, 0, time.UTC) @@ -94,9 +94,10 @@ func TestListNodeMetricHourly_FallsBackWhenRollupOnlyCoversRecentHours(t *testin }}}, nil } if strings.Contains(query, nodeMetricSnapshotTableName()) { - return &mockRows{data: [][]any{{ - rawHour, 12.0, 34.0, int64(5), int64(6), int64(7), int64(8), uint64(1), - }}}, nil + return &mockRows{data: [][]any{ + {rawHour, 12.0, 34.0, int64(5), int64(6), int64(7), int64(8), uint64(1)}, + {rollupHour, 50.0, 50.0, int64(9), int64(9), int64(9), int64(9), uint64(1)}, + }}, nil } return &mockRows{}, nil }, @@ -106,13 +107,33 @@ func TestListNodeMetricHourly_FallsBackWhenRollupOnlyCoversRecentHours(t *testin rows, err := ListNodeMetricHourly(ctx, NodeObservabilityFilter{Since: since}) require.NoError(t, err) - require.Len(t, rows, 1) - assert.Equal(t, 12.0, rows[0].AverageCPUUsagePercent) + require.Len(t, rows, 2) assert.Equal(t, rawHour, rows[0].Hour) + assert.Equal(t, 12.0, rows[0].AverageCPUUsagePercent) + // Overlapping hour prefers rollup (99) over raw (50). + assert.Equal(t, rollupHour, rows[1].Hour) + assert.Equal(t, 99.0, rows[1].AverageCPUUsagePercent) require.GreaterOrEqual(t, len(mock.queries), 2) assert.Contains(t, mock.queries[1], "lagInFrame") } +func TestMergeNodeMetricHourlyPreferRollup(t *testing.T) { + h1 := time.Date(2026, 7, 10, 10, 0, 0, 0, time.UTC) + h2 := time.Date(2026, 7, 10, 11, 0, 0, 0, time.UTC) + merged := mergeNodeMetricHourlyPreferRollup( + []NodeMetricHourly{{Hour: h2, AverageCPUUsagePercent: 80}}, + []NodeMetricHourly{ + {Hour: h1, AverageCPUUsagePercent: 10}, + {Hour: h2, AverageCPUUsagePercent: 20}, + }, + ) + require.Len(t, merged, 2) + assert.Equal(t, h1, merged[0].Hour) + assert.Equal(t, 10.0, merged[0].AverageCPUUsagePercent) + assert.Equal(t, h2, merged[1].Hour) + assert.Equal(t, 80.0, merged[1].AverageCPUUsagePercent) +} + func TestHourlyRollupCoversWindow(t *testing.T) { since := time.Date(2026, 7, 10, 0, 0, 0, 0, time.UTC) assert.True(t, hourlyRollupCoversWindow(since, since))