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.
This commit is contained in:
ryan
2026-07-10 11:24:31 +08:00
parent 4b11279662
commit da1dd92404
4 changed files with 191 additions and 22 deletions
+1 -1
View File
@@ -20,7 +20,7 @@ sidebar: false
### 修复
- 修复小时预聚合表仅有迁移后少量数据时仍优先于全量 raw 聚合、导致 24 小时容量/网络/磁盘趋势只显示最近时段的问题:rollup 未覆盖查询窗口起点时回退到原始快照小时聚合。
- 修复小时预聚合仅有迁移后少量数据时 24 小时容量/网络/磁盘趋势残缺的问题:读路径按小时 merge(rollup 覆盖不足时用 raw 补洞;窗口完整时仅走 rollup);并增加历史 backfill 迁移避免长期依赖 raw。
## [v3.1.2] - 2026-07-10
@@ -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;
@@ -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) {
@@ -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))