From 08d28c2c8ebfb05a3ead1630fa6fd887121b0f8e Mon Sep 17 00:00:00 2001 From: ryan Date: Sun, 9 Aug 2026 09:33:59 +0800 Subject: [PATCH] feat(log): drop empty old-month PG partitions during cleanup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 系统垃圾清理任务删除过期日志后,顺带清理旧月份空分区表: PostgreSQL 按月分区的访问日志表(节点/用户)在数据删除后若该月 分区已无数据,则自动删除对应分区表,避免历史分区表无限累积。 - 仅删除「当前月之前」且为空的月份分区,当月/未来月及仍有数据的分区保留 - ClickHouse/SQLite 为 no-op(CH 分区随数据删除自动消失) - 修复既有集成测试 pg_inherits 查询(inhrelid → inhparent) - 新增单元测试与 PG 集成测试 --- docs/changelog/index.md | 1 + internal/repository/logstore/cleanup.go | 40 ++++++- internal/repository/logstore/cleanup_test.go | 56 ++++++++++ .../repository/logstore/clickhouse_store.go | 5 + internal/repository/logstore/logstore.go | 3 + .../postgres_partition_integration_test.go | 101 +++++++++++++++++- .../repository/logstore/postgres_store.go | 33 ++++++ 7 files changed, 237 insertions(+), 2 deletions(-) diff --git a/docs/changelog/index.md b/docs/changelog/index.md index c2bebd39..730b8615 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -32,6 +32,7 @@ sidebar: false - 服务工作者(SW)注入挑战页改为前台无感知:不再显示「加载中…」文案,页面空白,仅通过浏览器控制台输出 `[sw-challenge]` 调试信息(注册成功/失败/回退重定向),注入过程不打扰访客。 - 性能指标(CPU/内存/磁盘/网络)与访问日志的保留时长解耦:新增三库共用的 `metric_retention_days` 配置(默认 3 天),系统垃圾清理每日任务按独立短留存清理指标快照;访问日志仍按 `log_retention_days_*` 清理。 - ClickHouse 改为默认关闭:`clickhouse.enabled` 缺省或为 `false` 时不启用(此前会被强制置为 `true`),日志/指标由 PostgreSQL/SQLite 主库承担;显式 `true` 或设置 `CLICKHOUSE_HOST` / `CLICKHOUSE_ENABLED=true` 时启用。 +- 系统垃圾清理任务在删除过期日志后,会同时清理旧月份空分区表:PostgreSQL 按月分区的访问日志表(节点/用户)在数据删除后若该月分区已无数据,则自动删除对应分区表,避免历史分区表无限累积;仅删除「当前月之前」且为空的月份分区,当月/未来月及仍有数据的分区保留。 ### 修复 diff --git a/internal/repository/logstore/cleanup.go b/internal/repository/logstore/cleanup.go index 91c617de..d74e2b52 100644 --- a/internal/repository/logstore/cleanup.go +++ b/internal/repository/logstore/cleanup.go @@ -8,6 +8,7 @@ import ( "errors" "fmt" "strconv" + "strings" "time" "github.com/Rain-kl/Wavelet/internal/model" @@ -37,6 +38,9 @@ const defaultMetricRetentionDays = 3 // partitionLeadMonths 清理时确保「当前月 + 未来 2 个月」分区持续存在。 const partitionLeadMonths = 2 +// accessLogPartitionTables 按月分区的访问日志表(分区预建/空分区清理共用)。 +var accessLogPartitionTables = []string{"of_node_access_logs", "w_user_access_logs"} + // retentionDaysForDatabase 按给定日志库读取保留天数(默认 90)。 func retentionDaysForDatabase(ctx context.Context, dbName string) int { key := model.ConfigKeyLogRetentionDaysPostgres @@ -108,6 +112,13 @@ func CleanupExpired(ctx context.Context) (*CleanupSummary, error) { }, summary); err != nil { return nil, err } + + // 过期数据删除后清理旧月份空分区表,避免分区表无限累积; + // 仅删「当前月之前」且无数据的分区(best-effort,失败不阻断数据保留清理)。 + if err := s.AccessLogs.DropEmptyPartitions(ctx, now); err != nil { + logger.WarnF(ctx, "drop empty log partitions failed: %v", err) + } + if err := cleanupTable("metric_snapshots", func() (int64, error) { return s.Observability.DeleteMetricSnapshotsBefore(ctx, metricCutoff) }, summary); err != nil { @@ -153,7 +164,7 @@ func partitionStatementsRange(from, to time.Time) []string { suffix := start.Format("200601") fromDay := start.Format("2006-01-02") toDay := monthEnd.Format("2006-01-02") - for _, table := range []string{"of_node_access_logs", "w_user_access_logs"} { + for _, table := range accessLogPartitionTables { out = append(out, fmt.Sprintf( "CREATE TABLE IF NOT EXISTS %s_%s PARTITION OF %s FOR VALUES FROM ('%s') TO ('%s')", table, suffix, table, fromDay, toDay)) @@ -161,3 +172,30 @@ func partitionStatementsRange(from, to time.Time) []string { } return out } + +// partitionNameMonth 解析按月分区表名 _YYYYMM 的所属月份;命名不匹配返回 (零值, false)。 +func partitionNameMonth(table, name string) (time.Time, bool) { + suffix, ok := strings.CutPrefix(name, table+"_") + if !ok || len(suffix) != 6 { + return time.Time{}, false + } + m, err := time.Parse("200601", suffix) + if err != nil { + return time.Time{}, false + } + return m, true +} + +// dropEligiblePartitionNames 返回 before 月份之前、命名合法的分区表名(是否为空由调用方校验)。 +func dropEligiblePartitionNames(table string, names []string, before time.Time) []string { + beforeMonth := time.Date(before.Year(), before.Month(), 1, 0, 0, 0, 0, time.UTC) + out := make([]string, 0, len(names)) + for _, name := range names { + month, ok := partitionNameMonth(table, name) + if !ok || !month.Before(beforeMonth) { + continue // 非法命名或当月/未来月分区,必须保留 + } + out = append(out, name) + } + return out +} diff --git a/internal/repository/logstore/cleanup_test.go b/internal/repository/logstore/cleanup_test.go index e1e20e3f..c2dab959 100644 --- a/internal/repository/logstore/cleanup_test.go +++ b/internal/repository/logstore/cleanup_test.go @@ -313,3 +313,59 @@ func hasAnySuffix(stmt string, suffixes []string) bool { } return false } + +// TestPartitionNameMonth 覆盖按月分区表名解析:合法命名返回所属月份,非法/其它表前缀返回 false。 +func TestPartitionNameMonth(t *testing.T) { + cases := []struct { + table string + name string + want string // 期望 "YYYY-MM";空串表示应解析失败 + }{ + {"of_node_access_logs", "of_node_access_logs_202608", "2026-08"}, + {"w_user_access_logs", "w_user_access_logs_202612", "2026-12"}, + {"of_node_access_logs", "w_user_access_logs_202608", ""}, // 其它表前缀 + {"of_node_access_logs", "of_node_access_logs_20268", ""}, // 位数不足 + {"of_node_access_logs", "of_node_access_logs_202613", ""}, // 非法月份 + {"of_node_access_logs", "of_node_access_logs_default", ""}, // 非数字后缀 + } + for _, c := range cases { + got, ok := partitionNameMonth(c.table, c.name) + if c.want == "" { + if ok { + t.Fatalf("partitionNameMonth(%q, %q) ok = true, want false", c.table, c.name) + } + continue + } + if !ok || got.Format("2006-01") != c.want { + t.Fatalf("partitionNameMonth(%q, %q) = %v, want %s", c.table, c.name, got, c.want) + } + } +} + +// TestDropEligiblePartitionNames 覆盖空分区清理筛选:只保留 before 月份之前、命名合法的分区。 +func TestDropEligiblePartitionNames(t *testing.T) { + before := time.Date(2026, 10, 15, 0, 0, 0, 0, time.UTC) + names := []string{ + "of_node_access_logs_202608", + "of_node_access_logs_202609", + "of_node_access_logs_202610", // 当月:保留 + "of_node_access_logs_202611", // 未来:保留 + "of_node_access_logs_default", // 非法命名:忽略 + } + got := dropEligiblePartitionNames("of_node_access_logs", names, before) + want := []string{"of_node_access_logs_202608", "of_node_access_logs_202609"} + if len(got) != len(want) { + t.Fatalf("eligible = %v, want %v", got, want) + } + for i, w := range want { + if got[i] != w { + t.Fatalf("eligible[%d] = %q, want %q", i, got[i], w) + } + } + + // 月初边界:before 恰为当月 1 日 0 点,当月分区仍保留。 + first := time.Date(2026, 10, 1, 0, 0, 0, 0, time.UTC) + if got := dropEligiblePartitionNames("of_node_access_logs", []string{"of_node_access_logs_202610"}, first); len(got) != 0 { + t.Fatalf("eligible at month boundary = %v, want empty", got) + } +} diff --git a/internal/repository/logstore/clickhouse_store.go b/internal/repository/logstore/clickhouse_store.go index 63eee7fe..62794ce6 100644 --- a/internal/repository/logstore/clickhouse_store.go +++ b/internal/repository/logstore/clickhouse_store.go @@ -434,6 +434,11 @@ func (s *clickhouseLogStore) EnsurePartitions(_ context.Context, _, _ time.Time) return nil } +// DropEmptyPartitions 是 CH 分支 no-op(CH 分区随数据删除自动消失,无独立分区表)。 +func (s *clickhouseLogStore) DropEmptyPartitions(_ context.Context, _ time.Time) error { + return nil +} + // chMigrationRange 查询 CH 表时间列 MIN/MAX;空表(NULL)返回零值。 func chMigrationRange(ctx context.Context, table, column string) (time.Time, time.Time, error) { if err := chConnErr(); err != nil { diff --git a/internal/repository/logstore/logstore.go b/internal/repository/logstore/logstore.go index cf5868da..72369dba 100644 --- a/internal/repository/logstore/logstore.go +++ b/internal/repository/logstore/logstore.go @@ -48,6 +48,9 @@ type AccessLogStore interface { // EnsurePartitions 幂等预建 PG 分区(按月),覆盖 [from, to] 月份;CH/SQLite 为 no-op。 // 目标为 PG 的迁移在复制前调用,避免历史数据写入报 "no partition of relation found"。 EnsurePartitions(ctx context.Context, from, to time.Time) error + // DropEmptyPartitions 幂等清理 PG 空分区表:删除 before 月份之前、且无任何数据的按月分区; + // CH/SQLite 为 no-op(CH 分区随数据删除自动消失、SQLite 无分区)。 + DropEmptyPartitions(ctx context.Context, before time.Time) error } // ObservabilityStore 可观测 4 表(metric snapshots / edge health / frps / frpc)。 diff --git a/internal/repository/logstore/postgres_partition_integration_test.go b/internal/repository/logstore/postgres_partition_integration_test.go index 11acce9e..16136ea1 100644 --- a/internal/repository/logstore/postgres_partition_integration_test.go +++ b/internal/repository/logstore/postgres_partition_integration_test.go @@ -86,7 +86,7 @@ func TestEnsurePartitionsPostgresInsertAcrossMonths(t *testing.T) { var partitionCount int64 if err := gdb.Raw( - "SELECT count(*) FROM pg_inherits WHERE inhrelid = to_regclass('of_node_access_logs')", + "SELECT count(*) FROM pg_inherits WHERE inhparent = to_regclass('of_node_access_logs')", ).Scan(&partitionCount).Error; err != nil { t.Fatalf("count partitions: %v", err) } @@ -146,6 +146,105 @@ func TestEnsurePartitionsPostgresInsertAcrossMonths(t *testing.T) { } } +// TestDropEmptyPartitionsPostgres 需要 TEST_POSTGRES_DSN(未设置时跳过): +// 验证空分区清理只删除 before 月份之前且无数据的分区:空旧月删除、有数据旧月保留、 +// 当月/未来月保留;用户访问日志分区同步清理。 +func TestDropEmptyPartitionsPostgres(t *testing.T) { + dsn := strings.TrimSpace(os.Getenv("TEST_POSTGRES_DSN")) + if dsn == "" { + t.Skip("TEST_POSTGRES_DSN is not set") + } + + gdb, err := gorm.Open(postgres.Open(dsn), &gorm.Config{ + DisableForeignKeyConstraintWhenMigrating: true, + Logger: logger.Default.LogMode(logger.Silent), + }) + if err != nil { + t.Fatalf("open postgres: %v", err) + } + sqlDB, err := gdb.DB() + if err != nil { + t.Fatalf("sql db: %v", err) + } + sqlDB.SetMaxOpenConns(1) + + schema := fmt.Sprintf("logstore_drop_partition_%d", time.Now().UnixNano()) + if !regexp.MustCompile(`^[a-z0-9_]+$`).MatchString(schema) { + t.Fatalf("invalid schema: %s", schema) + } + if err := gdb.Exec(`CREATE SCHEMA "` + schema + `"`).Error; err != nil { + t.Fatalf("create schema: %v", err) + } + if err := gdb.Exec(`SET search_path TO "` + schema + `"`).Error; err != nil { + t.Fatalf("set search_path: %v", err) + } + t.Cleanup(func() { + _ = gdb.Exec("SET search_path TO public").Error + _ = gdb.Exec(`DROP SCHEMA IF EXISTS "` + schema + `" CASCADE`).Error + _ = sqlDB.Close() + }) + + for _, ddl := range []string{postgresNodeAccessLogsDDL, postgresUserAccessLogsDDL} { + if err := gdb.Exec(ddl).Error; err != nil { + t.Fatalf("create partitioned table: %v", err) + } + } + + ctx := context.Background() + store := newGormStore(gdb) + ua := newUserAccessLogGormStore(gdb) + + // 预建 202601..202603 分区,仅 202602 有数据(节点+用户各 1 条),202601/202603 为空。 + if err := store.EnsurePartitions(ctx, + time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC), + time.Date(2026, 3, 20, 0, 0, 0, 0, time.UTC)); err != nil { + t.Fatalf("EnsurePartitions: %v", err) + } + if err := store.BatchInsertNodeAccessLogs(ctx, []analyticsmodel.NodeAccessLog{ + {ID: 1, NodeID: "n1", LoggedAt: time.Date(2026, 2, 10, 0, 0, 0, 0, time.UTC), RemoteAddr: "1.1.1.1"}, + }); err != nil { + t.Fatalf("insert node access log: %v", err) + } + if err := ua.BatchInsert(ctx, []analyticsmodel.UserAccessLog{ + {ID: 1, UserID: 101, Path: "/a", CreatedAt: time.Date(2026, 2, 11, 0, 0, 0, 0, time.UTC)}, + }); err != nil { + t.Fatalf("insert user access log: %v", err) + } + + // before=2026-03:202601(空)应删,202602(有数据)与 202603(当月)保留。 + if err := store.DropEmptyPartitions(ctx, time.Date(2026, 3, 15, 0, 0, 0, 0, time.UTC)); err != nil { + t.Fatalf("DropEmptyPartitions: %v", err) + } + + assertPartitions := func(parent string, want int64) { + t.Helper() + var n int64 + if err := gdb.Raw( + "SELECT count(*) FROM pg_inherits WHERE inhparent = to_regclass(?)", + parent, + ).Scan(&n).Error; err != nil { + t.Fatalf("count partitions of %s: %v", parent, err) + } + if n != want { + t.Fatalf("%s partitions = %d, want %d", parent, n, want) + } + } + assertPartitions("of_node_access_logs", 2) + assertPartitions("w_user_access_logs", 2) + + // 数据未受影响。 + var nodeCount, userCount int64 + if err := gdb.Model(&analyticsmodel.NodeAccessLog{}).Count(&nodeCount).Error; err != nil { + t.Fatalf("count node access logs: %v", err) + } + if err := gdb.Model(&analyticsmodel.UserAccessLog{}).Count(&userCount).Error; err != nil { + t.Fatalf("count user access logs: %v", err) + } + if nodeCount != 1 || userCount != 1 { + t.Fatalf("data counts = (%d, %d), want (1, 1)", nodeCount, userCount) + } +} + // postgresNodeAccessLogsDDL 与 goose/postgres/202608080001_create_log_tables.sql 对齐。 const postgresNodeAccessLogsDDL = ` CREATE TABLE IF NOT EXISTS of_node_access_logs ( diff --git a/internal/repository/logstore/postgres_store.go b/internal/repository/logstore/postgres_store.go index ca048b77..a4811d07 100644 --- a/internal/repository/logstore/postgres_store.go +++ b/internal/repository/logstore/postgres_store.go @@ -668,6 +668,39 @@ func (s *gormLogStore) EnsurePartitions(ctx context.Context, from, to time.Time) return nil } +// DropEmptyPartitions 幂等清理 PG 空分区表:仅删除 before 月份之前、且无任何数据的分区 +// (of_node_access_logs / w_user_access_logs 按月分区);非 PG 方言为 no-op(SQLite 无分区、CH 不走本实现)。 +func (s *gormLogStore) DropEmptyPartitions(ctx context.Context, before time.Time) error { + if !isPostgresDialect(s.db) { + return nil + } + for _, table := range accessLogPartitionTables { + var names []string + if err := s.db.WithContext(ctx).Raw(` +SELECT c.relname +FROM pg_inherits i +JOIN pg_class c ON c.oid = i.inhrelid +JOIN pg_class p ON p.oid = i.inhparent +JOIN pg_namespace n ON n.oid = p.relnamespace AND n.nspname = current_schema() +WHERE p.relname = ?`, table).Scan(&names).Error; err != nil { + return fmt.Errorf("list partitions of %s: %w", table, err) + } + for _, name := range dropEligiblePartitionNames(table, names, before) { + var one int + if err := s.db.WithContext(ctx).Raw("SELECT 1 FROM " + name + " LIMIT 1").Scan(&one).Error; err != nil { + return fmt.Errorf("check partition %s empty: %w", name, err) + } + if one == 1 { + continue // 仍有数据,保留 + } + if err := s.db.WithContext(ctx).Exec("DROP TABLE IF EXISTS " + name).Error; err != nil { + return fmt.Errorf("drop empty partition %s: %w", name, err) + } + } + } + return nil +} + // gormMigrationRange 按时间列 ORDER BY ± LIMIT 1 取首尾记录(经 GORM schema 扫描, // 避免 SQLite 时间存文本导致 MIN/MAX 原始 Scan 失败);空表返回两个零值。 func gormMigrationRange[T any](