From caf2ffcff4c3f27491702d917b39f6a498eb265c Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 2 Jul 2026 15:24:22 +0800 Subject: [PATCH] perf(clickhouse): P1 TTL migrations, ORDER BY tune, remote_addr normalization --- docs/changelog/index.md | 7 +++ docs/plan/clickhouse-cpu-optimization.md | 47 +++++++++++++++++++ .../apps/openflare/tasks/database_cleanup.go | 24 ++++++++++ .../openflare/tasks/database_cleanup_test.go | 2 +- ...202607020001_optimize_analytics_tables.sql | 23 +++++++++ internal/model/openflare_observability.go | 30 ++++++++++++ .../analytics/node_access_log_writer.go | 3 +- 7 files changed, 134 insertions(+), 2 deletions(-) create mode 100644 docs/plan/clickhouse-cpu-optimization.md create mode 100644 internal/db/migrator/goose/clickhouse/202607020001_optimize_analytics_tables.sql diff --git a/docs/changelog/index.md b/docs/changelog/index.md index fddb3d81..3eabfecf 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -20,6 +20,13 @@ sidebar: false ### 修改 +- ClickHouse 写入路径优化:移除 Agent 心跳路径中的同步 `ALTER DELETE` 保留清理;`batchwriter` 新增 `MinBatchSize` 抑制过小批次定时 flush;可观测 writer 批次提升至 500、flush 间隔 5s,并为 OpenResty/FRPS/FRPC 补全去重。 +- ClickHouse 客户端启用 `async_insert` 异步写入缓冲,并调高 `block_buffer_size` 与连接池默认值,降低小 part 与连接争用。 +- Dashboard 与节点可观测 API 消除无 `LIMIT` 全表扫描、增加短 TTL 内存缓存,前端轮询间隔分别调整为 60s/30s。 +- 访问日志与 WAF IP 组同步改为 ClickHouse 侧聚合与 SQL 分页,默认查询窗口限制为近 7 天,浏览器分布查询增加 Top 100 限制。 +- ClickHouse 分析表新增 TTL 自动过期策略:`w_user_access_logs` 180 天、`of_node_access_logs` 90 天,其余节点观测与聚合表 30 天;并优化 `of_node_access_logs` 的 `ORDER BY` 以匹配常见查询过滤模式。 +- 节点访问日志写入 ClickHouse 时对 `remote_addr` 执行 `TrimSpace` 规范化,避免首尾空白影响 IP 汇总统计。 +- 数据库自动清理任务新增 OpenResty、FRPS、FRPC 观测表清理目标。 - Docker 部署为 ClickHouse 服务增加 `nofile` ulimits 与 `docker/clickhouse/config.d/performance.xml` 性能配置挂载,限制 `max_concurrent_queries`、`background_pool_size` 与 `background_merges_mutations_concurrency_ratio`,降低高负载下的合并与查询争用。 - 审计访问日志写入 ClickHouse 时仅保留安全相关请求头(Authorization、Cookie、X-Forwarded-For、X-Real-IP、User-Agent、Content-Type),敏感头字段以 SHA-256 摘要脱敏,并将序列化后的 headers 载荷上限收紧至 2KB,减小 `w_user_access_logs` 行宽与 merge CPU 开销。 - 隐藏侧边栏“文档库”分组中的“规范示例”与“接口文档”,并将“使用文档”及其他相关页面的文档链接统一跳转至外部文档 https://open-flare.pages.dev/ diff --git a/docs/plan/clickhouse-cpu-optimization.md b/docs/plan/clickhouse-cpu-optimization.md new file mode 100644 index 00000000..37274b16 --- /dev/null +++ b/docs/plan/clickhouse-cpu-optimization.md @@ -0,0 +1,47 @@ +# ClickHouse CPU 性能优化计划 + +> PLAN_ID: `63ba981b` +> 状态: 执行中 +> 目标: 完成 P0–P2 优化,降低 ClickHouse CPU 占用 + +## 背景 + +ClickHouse CPU 偏高由写入侧(小 part 频繁 flush、心跳同步 DELETE mutation)与查询侧(无 LIMIT 全表扫、高频轮询、WAF 全量拉日志)叠加导致。 + +## PR Plan + +### PR 1: 写入路径 P0 优化 + +- **Description:** 移除心跳路径同步 `ALTER DELETE`;为 `batchwriter` 增加 `MinBatchSize`;调大可观测 writer 批次与 flush 间隔;为 openresty/frps/frpc 补全去重。 +- **Files/components affected:** `internal/apps/openflare/agent/observability.go`, `internal/db/batchwriter/`, `internal/apps/openflare/chwriter/`, `internal/db/batchwriter/*_test.go` +- **Dependencies:** None + +### PR 2: ClickHouse 客户端与配置 P1 + +- **Description:** 启用 `async_insert` 等写入优化 settings;提高 `block_buffer_size` 默认值;更新 `config.example.yaml` 与配置模型注释。 +- **Files/components affected:** `internal/db/clickhouse.go`, `internal/config/model.go`, `internal/config/config.go`, `config.example.yaml` +- **Dependencies:** None + +### PR 3: Dashboard 与可观测查询 P0 + +- **Description:** 消除 `limit=0` 无界查询;复用已有限制数据构建趋势;增加服务端短 TTL 缓存;降低前端轮询频率。 +- **Files/components affected:** `internal/apps/openflare/dashboard/logics.go`, `internal/apps/openflare/observability/node_logics.go`, `frontend/app/(main)/page.tsx`, `frontend/app/(main)/nodes/components/node-observability.tsx` +- **Dependencies:** None + +### PR 4: 访问日志与 WAF 查询 P0/P1 + +- **Description:** WAF IP 同步改为 ClickHouse 侧聚合;IP 汇总与折叠日志 SQL 分页;消除 count 重复全量扫描;列表 API 强制默认时间窗口。 +- **Files/components affected:** `internal/apps/openflare/waf/ip_group_sync.go`, `internal/repository/analytics/node_access_log_stats.go`, `internal/model/openflare_access_log.go`, `internal/apps/openflare/observability/access_log_logics.go`, `internal/repository/analytics/access_log_stats.go` +- **Dependencies:** None + +### PR 5: ClickHouse DDL 与数据规范化 P1 + +- **Description:** 为 7 张分析表添加 TTL;收窄 `of_node_access_logs` ORDER BY;插入时规范化 `remote_addr`(去 trim 查询);将可观测 obs 三表纳入自动清理。 +- **Files/components affected:** `internal/db/migrator/goose/clickhouse/`, `internal/repository/analytics/node_access_log_writer.go`, `internal/apps/openflare/tasks/database_cleanup.go`, `internal/model/analytics/` +- **Dependencies:** PR 1 + +### PR 6: 基础设施与审计减负 P2 + +- **Description:** Docker ClickHouse 服务端基础调优;审计日志 headers 截断/精简;更新 changelog。 +- **Files/components affected:** `docker-compose.yaml`, `docker/clickhouse/` (if needed), `internal/apps/risk_control/middleware.go`, `docs/changelog/index.md` +- **Dependencies:** None \ No newline at end of file diff --git a/internal/apps/openflare/tasks/database_cleanup.go b/internal/apps/openflare/tasks/database_cleanup.go index 96aea940..4d70e181 100644 --- a/internal/apps/openflare/tasks/database_cleanup.go +++ b/internal/apps/openflare/tasks/database_cleanup.go @@ -21,12 +21,21 @@ const ( DatabaseCleanupTargetMetricSnapshots = "node_metric_snapshots" // DatabaseCleanupTargetRequestReports is the API cleanup target for request reports. DatabaseCleanupTargetRequestReports = "node_request_reports" + // DatabaseCleanupTargetObsOpenresty is the API cleanup target for OpenResty observations. + DatabaseCleanupTargetObsOpenresty = "node_obs_openresty" + // DatabaseCleanupTargetObsFrps is the API cleanup target for FRPS observations. + DatabaseCleanupTargetObsFrps = "node_obs_frps" + // DatabaseCleanupTargetObsFrpc is the API cleanup target for FRPC observations. + DatabaseCleanupTargetObsFrpc = "node_obs_frpc" ) var databaseCleanupTargets = map[string]string{ DatabaseCleanupTargetAccessLogs: "访问日志", DatabaseCleanupTargetMetricSnapshots: "性能快照", DatabaseCleanupTargetRequestReports: "请求聚合", + DatabaseCleanupTargetObsOpenresty: "OpenResty 观测", + DatabaseCleanupTargetObsFrps: "FRPS 观测", + DatabaseCleanupTargetObsFrpc: "FRPC 观测", } // DatabaseCleanupInput describes a manual observability cleanup request. @@ -111,6 +120,9 @@ func RunDatabaseAutoCleanupOnce(ctx context.Context, now time.Time) (*DatabaseAu DatabaseCleanupTargetAccessLogs, DatabaseCleanupTargetMetricSnapshots, DatabaseCleanupTargetRequestReports, + DatabaseCleanupTargetObsOpenresty, + DatabaseCleanupTargetObsFrps, + DatabaseCleanupTargetObsFrpc, } { result, err := CleanupDatabaseObservability(ctx, DatabaseCleanupInput{ Target: target, @@ -137,6 +149,12 @@ func deleteAllObservabilityRows(ctx context.Context, target string) (int64, erro return model.DeleteAllOpenFlareMetricSnapshots(ctx) case DatabaseCleanupTargetRequestReports: return model.DeleteAllOpenFlareRequestReports(ctx) + case DatabaseCleanupTargetObsOpenresty: + return model.DeleteAllOpenFlareNodeObservationOpenresty(ctx) + case DatabaseCleanupTargetObsFrps: + return model.DeleteAllOpenFlareNodeObservationFrps(ctx) + case DatabaseCleanupTargetObsFrpc: + return model.DeleteAllOpenFlareNodeObservationFrpc(ctx) default: return 0, errors.New("unsupported cleanup target") } @@ -150,6 +168,12 @@ func deleteObservabilityRowsBefore(ctx context.Context, target string, cutoff ti return model.DeleteOpenFlareMetricSnapshotsBefore(ctx, cutoff) case DatabaseCleanupTargetRequestReports: return model.DeleteOpenFlareRequestReportsBefore(ctx, cutoff) + case DatabaseCleanupTargetObsOpenresty: + return model.DeleteOpenFlareNodeObservationOpenrestyBefore(ctx, cutoff) + case DatabaseCleanupTargetObsFrps: + return model.DeleteOpenFlareNodeObservationFrpsBefore(ctx, cutoff) + case DatabaseCleanupTargetObsFrpc: + return model.DeleteOpenFlareNodeObservationFrpcBefore(ctx, cutoff) default: return 0, errors.New("unsupported cleanup target") } diff --git a/internal/apps/openflare/tasks/database_cleanup_test.go b/internal/apps/openflare/tasks/database_cleanup_test.go index 408b006c..cfe65daa 100644 --- a/internal/apps/openflare/tasks/database_cleanup_test.go +++ b/internal/apps/openflare/tasks/database_cleanup_test.go @@ -131,7 +131,7 @@ func TestRunDatabaseAutoCleanupOnceDeletesAllObservabilityTargets(t *testing.T) summary, err := RunDatabaseAutoCleanupOnce(ctx, now) require.NoError(t, err) require.NotNil(t, summary) - require.Len(t, summary.Results, 3) + require.Len(t, summary.Results, 6) accessLogs, err := model.ListOpenFlareAccessLogs(ctx, model.OpenFlareAccessLogQuery{Page: 0, PageSize: 10}) require.NoError(t, err) diff --git a/internal/db/migrator/goose/clickhouse/202607020001_optimize_analytics_tables.sql b/internal/db/migrator/goose/clickhouse/202607020001_optimize_analytics_tables.sql new file mode 100644 index 00000000..8881237e --- /dev/null +++ b/internal/db/migrator/goose/clickhouse/202607020001_optimize_analytics_tables.sql @@ -0,0 +1,23 @@ +-- +goose Up +-- Add TTL policies to analytics tables so ClickHouse can expire rows automatically. +ALTER TABLE w_user_access_logs MODIFY TTL created_at + INTERVAL 180 DAY; + +ALTER TABLE of_node_access_logs MODIFY TTL logged_at + INTERVAL 90 DAY; + +ALTER TABLE of_node_metric_snapshots MODIFY TTL captured_at + INTERVAL 30 DAY; + +ALTER TABLE of_node_request_reports MODIFY TTL window_ended_at + INTERVAL 30 DAY; + +ALTER TABLE of_node_obs_openresty MODIFY TTL captured_at + INTERVAL 30 DAY; + +ALTER TABLE of_node_obs_frps MODIFY TTL captured_at + INTERVAL 30 DAY; + +ALTER TABLE of_node_obs_frpc MODIFY TTL captured_at + INTERVAL 30 DAY; + +-- Narrow ORDER BY for node access logs to match common filter patterns. +-- Requires ClickHouse 24.10+ (MODIFY ORDER BY). On older versions this statement +-- may fail and require manual table recreation; TTL changes above are still safe. +ALTER TABLE of_node_access_logs MODIFY ORDER BY (node_id, logged_at, status_code); + +-- +goose Down +-- TTL and ORDER BY changes cannot be safely reversed without recreating tables. \ No newline at end of file diff --git a/internal/model/openflare_observability.go b/internal/model/openflare_observability.go index d47b0099..4872bd21 100644 --- a/internal/model/openflare_observability.go +++ b/internal/model/openflare_observability.go @@ -396,6 +396,36 @@ func DeleteAllOpenFlareRequestReports(ctx context.Context) (int64, error) { return currentObservabilityStore().DeleteAllRequestReports(ctx) } +// DeleteOpenFlareNodeObservationOpenrestyBefore deletes OpenResty observations captured before cutoff. +func DeleteOpenFlareNodeObservationOpenrestyBefore(ctx context.Context, cutoff time.Time) (int64, error) { + return currentObservabilityStore().DeleteNodeObservationOpenrestyBefore(ctx, cutoff) +} + +// DeleteAllOpenFlareNodeObservationOpenresty deletes all OpenResty observations. +func DeleteAllOpenFlareNodeObservationOpenresty(ctx context.Context) (int64, error) { + return currentObservabilityStore().DeleteAllNodeObservationOpenresty(ctx) +} + +// DeleteOpenFlareNodeObservationFrpsBefore deletes FRPS observations captured before cutoff. +func DeleteOpenFlareNodeObservationFrpsBefore(ctx context.Context, cutoff time.Time) (int64, error) { + return currentObservabilityStore().DeleteNodeObservationFrpsBefore(ctx, cutoff) +} + +// DeleteAllOpenFlareNodeObservationFrps deletes all FRPS observations. +func DeleteAllOpenFlareNodeObservationFrps(ctx context.Context) (int64, error) { + return currentObservabilityStore().DeleteAllNodeObservationFrps(ctx) +} + +// DeleteOpenFlareNodeObservationFrpcBefore deletes FRPC observations captured before cutoff. +func DeleteOpenFlareNodeObservationFrpcBefore(ctx context.Context, cutoff time.Time) (int64, error) { + return currentObservabilityStore().DeleteNodeObservationFrpcBefore(ctx, cutoff) +} + +// DeleteAllOpenFlareNodeObservationFrpc deletes all FRPC observations. +func DeleteAllOpenFlareNodeObservationFrpc(ctx context.Context) (int64, error) { + return currentObservabilityStore().DeleteAllNodeObservationFrpc(ctx) +} + // DeleteOpenFlareHealthEventsByNodeID deletes all health events for a node. func DeleteOpenFlareHealthEventsByNodeID(ctx context.Context, nodeID string) (int64, error) { conn := db.DB(ctx) diff --git a/internal/repository/analytics/node_access_log_writer.go b/internal/repository/analytics/node_access_log_writer.go index 58fe6ef8..88b02336 100644 --- a/internal/repository/analytics/node_access_log_writer.go +++ b/internal/repository/analytics/node_access_log_writer.go @@ -6,6 +6,7 @@ package analytics import ( "context" "fmt" + "strings" "time" "github.com/Rain-kl/Wavelet/internal/db" @@ -41,7 +42,7 @@ func BatchInsertNodeAccessLogs(ctx context.Context, logs []analyticsmodel.NodeAc id, logItem.NodeID, logItem.LoggedAt.UTC(), - logItem.RemoteAddr, + strings.TrimSpace(logItem.RemoteAddr), logItem.Region, logItem.Host, logItem.Path,