diff --git a/docs/changelog/index.md b/docs/changelog/index.md index d5feebc9..1c8a8551 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -18,6 +18,10 @@ sidebar: false ## [unreleased] +### 修复 + +- 修复小时预聚合表仅有迁移后少量数据时仍优先于全量 raw 聚合、导致 24 小时容量/网络/磁盘趋势只显示最近时段的问题:rollup 未覆盖查询窗口起点时回退到原始快照小时聚合。 + ## [v3.1.2] - 2026-07-10 ### 修复 diff --git a/internal/repository/analytics/node_observability.go b/internal/repository/analytics/node_observability.go index 8adc8c4b..467ad375 100644 --- a/internal/repository/analytics/node_observability.go +++ b/internal/repository/analytics/node_observability.go @@ -368,11 +368,31 @@ ORDER BY hour ASC`, nodeTrafficHourlyTableName, clause) return result, nil } +// 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. +const hourlyRollupMaxLead = 2 * time.Hour + +// hourlyRollupCoversWindow reports whether rollup coverage starts near the requested window. +// rows must be ordered by hour ascending. +func hourlyRollupCoversWindow(earliestHour time.Time, since time.Time) bool { + if since.IsZero() { + return true + } + sinceHour := since.UTC().Truncate(time.Hour) + earliest := earliestHour.UTC().Truncate(time.Hour) + return !earliest.After(sinceHour.Add(hourlyRollupMaxLead)) +} + // ListNodeMetricHourly returns hourly metric snapshot aggregates matching filter. -// Prefers of_node_metric_capacity_hourly rollups; falls back to raw lagInFrame on empty/error. +// Prefers of_node_metric_capacity_hourly when it covers the requested window; otherwise +// falls back to raw lagInFrame over of_node_metric_snapshots. func ListNodeMetricHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeMetricHourly, error) { if rows, err := listNodeMetricHourlyFromRollup(ctx, filter); err == nil && len(rows) > 0 { - return rows, nil + if hourlyRollupCoversWindow(rows[0].Hour, filter.Since) { + return rows, nil + } } return listNodeMetricHourlyFromRaw(ctx, filter) } @@ -492,10 +512,13 @@ func scanNodeMetricHourlyRows(rows driver.Rows) ([]NodeMetricHourly, error) { } // ListNodeOpenrestyHourly returns hourly OpenResty observation aggregates matching filter. -// Prefers of_node_openresty_hourly rollups; falls back to raw lagInFrame on empty/error. +// Prefers of_node_openresty_hourly when it covers the requested window; otherwise +// falls back to raw lagInFrame over of_node_obs_openresty. func ListNodeOpenrestyHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) { if rows, err := listNodeOpenrestyHourlyFromRollup(ctx, filter); err == nil && len(rows) > 0 { - return rows, nil + if hourlyRollupCoversWindow(rows[0].Hour, filter.Since) { + return rows, nil + } } return listNodeOpenrestyHourlyFromRaw(ctx, filter) } diff --git a/internal/repository/analytics/node_observability_latest_test.go b/internal/repository/analytics/node_observability_latest_test.go index 8478e4fb..03c89dc0 100644 --- a/internal/repository/analytics/node_observability_latest_test.go +++ b/internal/repository/analytics/node_observability_latest_test.go @@ -55,6 +55,7 @@ func TestListLatestNodeRequestReports_UsesLimit1ByNodeID(t *testing.T) { func TestListNodeMetricHourly_PrefersRollup(t *testing.T) { ctx := context.Background() hour := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + since := hour.Add(-1 * time.Hour) mock := &mockConn{ queryFn: func(_ context.Context, query string, _ ...any) (driver.Rows, error) { if strings.Contains(query, nodeMetricCapacityHourlyTableName()) { @@ -62,13 +63,13 @@ func TestListNodeMetricHourly_PrefersRollup(t *testing.T) { hour, 42.5, 60.0, int64(100), int64(200), int64(10), int64(20), uint64(2), }}}, nil } - return nil, errors.New("raw path should not be used when rollup has rows") + return nil, errors.New("raw path should not be used when rollup covers the window") }, } db.SetChConnForTest(mock) t.Cleanup(func() { db.SetChConnForTest(nil) }) - rows, err := ListNodeMetricHourly(ctx, NodeObservabilityFilter{}) + rows, err := ListNodeMetricHourly(ctx, NodeObservabilityFilter{Since: since}) require.NoError(t, err) require.Len(t, rows, 1) assert.Equal(t, 42.5, rows[0].AverageCPUUsagePercent) @@ -79,6 +80,47 @@ func TestListNodeMetricHourly_PrefersRollup(t *testing.T) { assert.Contains(t, mock.queries[0], nodeMetricCapacityHourlyTableName()) } +func TestListNodeMetricHourly_FallsBackWhenRollupOnlyCoversRecentHours(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) + rollupHour := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + rawHour := time.Date(2026, 7, 9, 15, 0, 0, 0, time.UTC) + mock := &mockConn{ + queryFn: func(_ context.Context, query string, _ ...any) (driver.Rows, error) { + if strings.Contains(query, nodeMetricCapacityHourlyTableName()) { + return &mockRows{data: [][]any{{ + rollupHour, 99.0, 99.0, int64(1), int64(1), int64(1), int64(1), uint64(1), + }}}, 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{}, nil + }, + } + db.SetChConnForTest(mock) + t.Cleanup(func() { db.SetChConnForTest(nil) }) + + rows, err := ListNodeMetricHourly(ctx, NodeObservabilityFilter{Since: since}) + require.NoError(t, err) + require.Len(t, rows, 1) + assert.Equal(t, 12.0, rows[0].AverageCPUUsagePercent) + assert.Equal(t, rawHour, rows[0].Hour) + require.GreaterOrEqual(t, len(mock.queries), 2) + assert.Contains(t, mock.queries[1], "lagInFrame") +} + +func TestHourlyRollupCoversWindow(t *testing.T) { + since := time.Date(2026, 7, 10, 0, 0, 0, 0, time.UTC) + assert.True(t, hourlyRollupCoversWindow(since, since)) + assert.True(t, hourlyRollupCoversWindow(since.Add(2*time.Hour), since)) + assert.False(t, hourlyRollupCoversWindow(since.Add(3*time.Hour), since)) + assert.True(t, hourlyRollupCoversWindow(time.Date(2026, 7, 11, 0, 0, 0, 0, time.UTC), time.Time{})) +} + func TestListNodeMetricHourly_FallsBackToRawOnRollupError(t *testing.T) { ctx := context.Background() hour := time.Date(2026, 7, 10, 13, 0, 0, 0, time.UTC) @@ -111,6 +153,7 @@ func TestListNodeMetricHourly_FallsBackToRawOnRollupError(t *testing.T) { func TestListNodeOpenrestyHourly_PrefersRollup(t *testing.T) { ctx := context.Background() hour := time.Date(2026, 7, 10, 14, 0, 0, 0, time.UTC) + since := hour.Add(-1 * time.Hour) mock := &mockConn{ queryFn: func(_ context.Context, query string, _ ...any) (driver.Rows, error) { if strings.Contains(query, nodeOpenrestyHourlyTableName()) { @@ -118,13 +161,13 @@ func TestListNodeOpenrestyHourly_PrefersRollup(t *testing.T) { hour, int64(50), int64(70), uint64(3), }}}, nil } - return nil, errors.New("raw path should not be used when rollup has rows") + return nil, errors.New("raw path should not be used when rollup covers the window") }, } db.SetChConnForTest(mock) t.Cleanup(func() { db.SetChConnForTest(nil) }) - rows, err := ListNodeOpenrestyHourly(ctx, NodeObservabilityFilter{}) + rows, err := ListNodeOpenrestyHourly(ctx, NodeObservabilityFilter{Since: since}) require.NoError(t, err) require.Len(t, rows, 1) assert.Equal(t, int64(50), rows[0].OpenrestyRxBytes) @@ -159,8 +202,3 @@ func TestListNodeOpenrestyHourly_FallsBackToRawOnEmptyRollup(t *testing.T) { assert.Contains(t, mock.queries[1], "lagInFrame") } -func TestNonNegativeCounterRange(t *testing.T) { - assert.Equal(t, int64(0), nonNegativeCounterRange(10, 10)) - assert.Equal(t, int64(0), nonNegativeCounterRange(5, 10)) - assert.Equal(t, int64(15), nonNegativeCounterRange(25, 10)) -} diff --git a/scripts/live_ch_smoke/main.go b/scripts/live_ch_smoke/main.go index ed32296e..8357f191 100644 --- a/scripts/live_ch_smoke/main.go +++ b/scripts/live_ch_smoke/main.go @@ -1,3 +1,7 @@ +// Package main is a manual smoke tool for ClickHouse app write path. +// Usage (from repo root, with config.yaml and Docker CH up): +// +// go run ./scripts/live_ch_smoke package main import ( @@ -11,13 +15,27 @@ import ( "github.com/Rain-kl/Wavelet/internal/model" ) +const ( + flushWaitTimeout = 45 * time.Second + listLimit = 5 + pollInterval = 2 * time.Second +) + func main() { - if !db.ChConnReady() { - fmt.Fprintln(os.Stderr, "ChConn not ready — check config.yaml clickhouse.enabled") + if err := run(); err != nil { + fmt.Fprintln(os.Stderr, err) os.Exit(1) } +} + +func run() error { + if !db.ChConnReady() { + return fmt.Errorf("ChConn not ready — check config.yaml clickhouse.enabled") + } ctx := context.Background() chwriter.Init(ctx) + defer func() { _ = chwriter.Stop(ctx) }() + now := time.Now().UTC() nodeID := "e2e-app-" + now.Format("150405") if err := model.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{ @@ -26,44 +44,49 @@ func main() { StorageUsedBytes: 456, StorageTotalBytes: 2000, DiskReadBytes: 11, DiskWriteBytes: 22, NetworkRxBytes: 33, NetworkTxBytes: 44, }); err != nil { - fmt.Fprintln(os.Stderr, "insert:", err) - os.Exit(1) + return fmt.Errorf("insert: %w", err) } fmt.Println("queued", nodeID) - deadline := time.Now().Add(45 * time.Second) + + if err := waitForSnapshot(ctx, nodeID, now); err != nil { + return err + } + if err := assertLatestIncludes(ctx, nodeID, now); err != nil { + return err + } + for _, s := range chwriter.WriterStats() { + fmt.Printf("writer %s running=%v depth=%d drops=%d flush_err=%d\n", + s.Name, s.Running, s.Depth, s.Drops, s.FlushErrors) + } + return nil +} + +func waitForSnapshot(ctx context.Context, nodeID string, now time.Time) error { + deadline := time.Now().Add(flushWaitTimeout) for time.Now().Before(deadline) { - rows, err := model.ListOpenFlareMetricSnapshotsSince(ctx, nodeID, now.Add(-time.Minute), 5) + rows, err := model.ListOpenFlareMetricSnapshotsSince(ctx, nodeID, now.Add(-time.Minute), listLimit) if err != nil { - fmt.Fprintln(os.Stderr, "list:", err) - os.Exit(1) + return fmt.Errorf("list: %w", err) } if len(rows) > 0 { fmt.Printf("OK flushed id=%d cpu=%.1f\n", rows[0].ID, rows[0].CPUUsagePercent) - latest, err := model.ListOpenFlareLatestMetricSnapshotsSince(ctx, "", now.Add(-time.Hour)) - if err != nil { - fmt.Fprintln(os.Stderr, "latest:", err) - os.Exit(1) - } - ok := false - for _, r := range latest { - if r != nil && r.NodeID == nodeID { - ok = true - } - } - if !ok { - fmt.Fprintln(os.Stderr, "FAIL latest-per-node missing node") - os.Exit(1) - } - fmt.Println("OK latest-per-node includes node") - for _, s := range chwriter.WriterStats() { - fmt.Printf("writer %s running=%v depth=%d drops=%d flush_err=%d\n", - s.Name, s.Running, s.Depth, s.Drops, s.FlushErrors) - } - _ = chwriter.Stop(ctx) - return + return nil } - time.Sleep(2 * time.Second) + time.Sleep(pollInterval) } - fmt.Fprintln(os.Stderr, "FAIL: not flushed within 45s") - os.Exit(1) + return fmt.Errorf("not flushed within timeout") +} + +func assertLatestIncludes(ctx context.Context, nodeID string, now time.Time) error { + latest, err := model.ListOpenFlareLatestMetricSnapshotsSince(ctx, "", now.Add(-time.Hour)) + if err != nil { + return fmt.Errorf("latest: %w", err) + } + for _, r := range latest { + if r != nil && r.NodeID == nodeID { + fmt.Println("OK latest-per-node includes node") + return nil + } + } + return fmt.Errorf("latest-per-node missing node %s", nodeID) }