mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 05:46:36 +08:00
chore: docs
This commit is contained in:
@@ -7,7 +7,7 @@ description: "Wavelet 项目专用:当新增或修改 ClickHouse 批量写入
|
||||
|
||||
开始前阅读根目录 `AGENTS.md`。ClickHouse 是辅助 OLAP 存储,**厌恶高频单条写入**(过多小 part);写入路径必须优先批量或异步聚合。
|
||||
|
||||
DDL 与表结构变更见 `database-migration` 技能;本技能只覆盖**运行时写入架构**。
|
||||
DDL 与表结构变更见 `database-migration` 技能。日志/分析用途表的判定、三库回落与切换见 `logstore` 技能。本技能只覆盖**运行时写入架构**。
|
||||
|
||||
## 分层职责
|
||||
|
||||
@@ -17,7 +17,7 @@ DDL 与表结构变更见 `database-migration` 技能;本技能只覆盖**运
|
||||
| 批量框架 | `internal/infra/persistence/batchwriter/` | 泛型队列 + 按条数/时间 flush + 非阻塞入队 + 优雅停机;**各业务域独立实例** |
|
||||
| Model | `internal/model/analytics/` | 列定义、`TableName()`、`BatchInsertSQL()`(及可选 `InsertColumns()`) |
|
||||
| Repository | `internal/repository/analytics/` | `BatchInsert*` / `BatchInsertNodeAccessLogs` 等;`PrepareBatch` + 多行 `Append` + 一次 `Send` |
|
||||
| Apps | `internal/apps/<domain>/` | 采集、入队、背压;`FlushFunc` 只调 repository,不写 SQL、不 `PrepareBatch` |
|
||||
| Apps | `internal/apps/<domain>/` | 采集、入队、背压;`FlushFunc` 只调 logstore / repository,不写 SQL、不 `PrepareBatch` |
|
||||
| 装配 | `internal/platform/bootstrap/bootstrap.go` | 进程启动时调用 `Writer.Start`;初始化时需调用 `lifecycle.OnShutdown` 挂载停机钩子 |
|
||||
| 生命周期 | `internal/platform/lifecycle/lifecycle.go` | 统一协调全局并发优雅停机,业务包无需在 `bootstrap.go` 中硬编码 `Stop` 逻辑 |
|
||||
|
||||
@@ -37,6 +37,7 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush
|
||||
|
||||
- `QueueSize`: 10_000
|
||||
- `MaxBatchSize`: 1_000
|
||||
- `MinBatchSize`: 50(未达阈值则跳过按时间 flush,除非设了 `MaxFlushWait`)
|
||||
- `FlushInterval`: 1s
|
||||
|
||||
各域可独立覆盖;可观测低频指标可用更小 `MaxBatchSize`(如 100)与更长 `FlushInterval`(如 2–5s),但**不要**退化为逐条 `Send`。
|
||||
@@ -49,7 +50,8 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush
|
||||
### FlushFunc 规范
|
||||
|
||||
- 签名:`func(ctx context.Context, items []T) error`
|
||||
- 内部调用 `internal/repository/analytics` 的 `BatchInsert*`(传入 `[]analyticsmodel.X`)
|
||||
- **日志/分析用途表**:`logstore.Active(ctx)` 再调对应 `BatchInsert*`。禁止 apps 直连 `analyticsrepo` 或 `db.ChConn`。
|
||||
- 仅 CH、无需主库回落的分析表:才直接调 `repository/analytics` 的 `BatchInsert*`。
|
||||
- 在 flush 边界记录一次错误日志,不要把 DB 驱动错误直接暴露给 HTTP 客户端
|
||||
- `Start` 使用 `context.WithoutCancel(parent)`,避免请求 ctx 取消中断后台 flush
|
||||
|
||||
@@ -57,11 +59,11 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush
|
||||
|
||||
每个业务域拥有自己的 `Writer`、配置与 `FlushFunc`:
|
||||
|
||||
| 域 | 表 | 现状 | 目标形态 |
|
||||
| :--- | :--- | :--- | :--- |
|
||||
| 管理端审计 | `w_user_access_logs` | `risk_control` → `batchwriter` + `analyticsrepo.BatchInsert` | 已接入 |
|
||||
| 边缘访问日志 | `of_node_access_logs` | `openflare/chwriter` 异步 flush | 已接入 |
|
||||
| 可观测时序 | `of_node_metric_snapshots` 等 5 表 | `openflare/chwriter` 五表独立 writer + 进程内短 TTL 去重 | 已接入 |
|
||||
| 域 | 表 | 写入路径 |
|
||||
| :--- | :--- | :--- |
|
||||
| 管理端审计 | `w_user_access_logs` | `risk_control` → `batchwriter` → `logstore.Active` |
|
||||
| 边缘访问日志 | `of_node_access_logs` | `openflare/chwriter` → `logstore.Active` |
|
||||
| 可观测时序 | `of_node_metric_snapshots` 等 | `openflare/chwriter` 分表 writer + 进程内短 TTL 去重 → `logstore.Active` |
|
||||
|
||||
**不要**把 audit、access log、observability 并入同一 channel。
|
||||
|
||||
@@ -73,8 +75,9 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush
|
||||
- `len(items)==0` 直接返回
|
||||
- `db.ChConn == nil` 返回明确错误
|
||||
- 一次 `PrepareBatch` → 循环 `Append` → 一次 `Send`
|
||||
4. **Writer 胶水**(`internal/apps/<domain>/` 或 `internal/repository/analytics/<domain>_writer.go`):
|
||||
4. **Writer 胶水**(`internal/apps/<domain>/`):
|
||||
- `New` + `Start`,并在初始化逻辑内通过 `lifecycle.OnShutdown("your_writer_name", Stop)` 注册停机回调
|
||||
- 日志表的 `FlushFunc` 调 `logstore.Active`(见 `logstore` skill)
|
||||
- 业务路径 `TryEnqueue`;HTTP 背压用 `IsFull()`
|
||||
5. **测试**:
|
||||
- repository:mock `ChConn` 验证 `BatchInsertSQL` 与 append 列数
|
||||
@@ -123,14 +126,9 @@ var globalChan chan any
|
||||
|
||||
```go
|
||||
// internal/platform/bootstrap/bootstrap.go(示意)
|
||||
var userAccessLogWriter *batchwriter.Writer[*analytics.UserAccessLog]
|
||||
|
||||
func RegisterAPI(ctx context.Context) {
|
||||
// ...
|
||||
if config.Config.ClickHouse.Enabled {
|
||||
initUserAccessLogWriter(ctx) // Start writer
|
||||
risk_control.BindWriter(userAccessLogWriter) // 或逐步替换 InitLogWriter
|
||||
}
|
||||
// 日志 writer 不依赖 clickhouse.enabled:flush 时由 logstore 选库
|
||||
risk_control.InitLogWriter(ctx)
|
||||
}
|
||||
```
|
||||
|
||||
@@ -149,7 +147,8 @@ make code-check
|
||||
- flush 按 `MaxBatchSize` 与 `FlushInterval` 触发
|
||||
- `Stop` 能 drain 队列内剩余项
|
||||
- repository 层无 goroutine、无 channel
|
||||
- `clickhouse.enabled: false` 时不 `Start` writer、不入队
|
||||
- 日志表:`clickhouse.enabled: false` 时 writer 仍 `Start`,flush 走主库 logstore
|
||||
- 仅 CH 的分析表:未启用 CH 时不要 `Start`、不要入队
|
||||
|
||||
## 相关文件速查
|
||||
|
||||
@@ -157,7 +156,8 @@ make code-check
|
||||
- 连接:`internal/infra/persistence/clickhouse.go`
|
||||
- 审计写入:`internal/apps/risk_control/logics.go`
|
||||
- OpenFlare 写入胶水:`internal/apps/openflare/chwriter/writer.go`
|
||||
- 节点访问日志 repository:`internal/repository/analytics/node_access_log_writer.go`
|
||||
- 可观测 repository:`internal/repository/analytics/node_observability_writer.go`
|
||||
- 日志抽象:`internal/repository/logstore`
|
||||
- 节点访问日志 CH 实现:`internal/repository/analytics/node_access_log_writer.go`
|
||||
- 可观测 CH 实现:`internal/repository/analytics/node_observability_writer.go`
|
||||
- 生命周期管理器:`internal/platform/lifecycle/lifecycle.go`
|
||||
- Bootstrap:`internal/platform/bootstrap/bootstrap.go`
|
||||
@@ -81,7 +81,7 @@ make code-check
|
||||
ClickHouse 是**辅助 OLAP 存储**,与 PostgreSQL/SQLite 主库**完全独立**的迁移与访问管线:
|
||||
|
||||
- 主库(PG/SQLite):业务事务数据、`goose_db_version`、双方言 SQL。
|
||||
- 分析库(ClickHouse):访问日志、统计聚合等分析型数据、`goose_clickhouse_version`、单方言 SQL。
|
||||
- 分析库(ClickHouse):分析型数据、`goose_clickhouse_version`、单方言 SQL。日志用途表还必须在主库建回落并走 `logstore`(见该 skill);CH 目录仍只放 CH DDL。
|
||||
|
||||
**不要**把 ClickHouse 表结构混入 PG/SQLite 迁移目录,也**不要**在 `support-files/`、`internal/apps/` 或 `internal/repository/` 中手写 DDL。
|
||||
|
||||
@@ -118,7 +118,7 @@ ClickHouse 是**辅助 OLAP 存储**,与 PostgreSQL/SQLite 主库**完全独
|
||||
1. **Model**:在 `internal/model/analytics/` 定义 struct,`gorm:"column:..."` 与 DDL 列名一一对应;实现 `TableName()`,批量写入表可提供 `InsertColumns()` / `BatchInsertSQL()`。
|
||||
2. **Goose SQL**:在 `internal/infra/persistence/migrator/goose/clickhouse/` 新增递增版本文件(格式同主库,如 `YYYYMMDDNNNN_create_xxx.sql`),编写 `-- +goose Up` / `-- +goose Down`。
|
||||
3. **Repository**:在 `internal/repository/analytics/` 实现 `BatchInsert*`(`db.ChConn` 一次 `PrepareBatch` + 多行 `Append` + 一次 `Send`)与查询(`db.ChDB`);连接未初始化时返回明确错误,**不要**在 handler 写 SQL,**不要**在 repository 内维护 channel/goroutine。
|
||||
4. **Apps**:在 `internal/apps/<domain>/` 编排采集与入队;高频写入通过 `internal/infra/persistence/batchwriter` 各域独立实例异步 flush(详见 `clickhouse-batchwriter` 技能),`FlushFunc` 只调 repository `BatchInsert*`;管理端统计 API 只读 repository,不触达 DDL。
|
||||
4. **Apps**:在 `internal/apps/<domain>/` 编排采集与入队;高频写入通过 `internal/infra/persistence/batchwriter` 各域独立实例异步 flush(详见 `clickhouse-batchwriter` 技能)。**日志/分析用途表**还要同时建 PG/SQLite 回落并接入 `logstore`(见 `logstore` 技能),`FlushFunc` 调 `logstore.Active` 而不是 `analyticsrepo`;普通业务分析表仍只读 repository。
|
||||
|
||||
### ClickHouse 验证
|
||||
|
||||
|
||||
@@ -0,0 +1,75 @@
|
||||
---
|
||||
name: "logstore"
|
||||
description: "OpenFlare / Wavelet:当新增或修改日志/分析用途表(节点访问日志、用户访问日志、可观测时序)、接入 internal/repository/logstore、切换日志主库、实现 PG/SQLite 回落,或判断一张表该走业务主库还是日志库时必须使用。"
|
||||
---
|
||||
|
||||
# 日志用途表开发
|
||||
|
||||
开始前阅读根目录 `AGENTS.md`。DDL 用 `database-migration`;高频写入队列用 `clickhouse-batchwriter`;切换任务用 `new-async-task`。本技能只回答:**这张表是不是日志表,以及如何接入可切换的日志主库。**
|
||||
|
||||
设计背景见 [日志存储解耦](../../../docs/design/logstore.md)。
|
||||
|
||||
## 先判定
|
||||
|
||||
日志表同时满足:
|
||||
|
||||
- 追加写入、几乎不更新单行
|
||||
- 按时间查询/聚合,允许按保留天数删除
|
||||
- 关闭 ClickHouse 后仍要能写、能查
|
||||
- 不参与网站/节点/证书等事务一致性
|
||||
|
||||
**不要**做成日志表:Zone、节点、配置版本、任务执行、上传元数据。这些走主库 `repository`。
|
||||
|
||||
当前日志域:
|
||||
|
||||
| 域 | 接口 | 表 |
|
||||
| :--- | :--- | :--- |
|
||||
| 节点访问日志 | `AccessLogStore` | `of_node_access_logs` |
|
||||
| 可观测 | `ObservabilityStore` | `of_node_metric_snapshots` / `of_node_edge_health` / `of_node_obs_frps` / `of_node_obs_frpc` |
|
||||
| 用户访问审计 | `UserAccessLogStore` | `w_user_access_logs` |
|
||||
|
||||
## 分层
|
||||
|
||||
| 层级 | 路径 | 职责 |
|
||||
| :--- | :--- | :--- |
|
||||
| 抽象 | `internal/repository/logstore` | 接口 + `Active`/`BuildForMigration`;apps **只**面向这里或 `repository` 门面 |
|
||||
| CH 实现 | `logstore/clickhouse_store.go` 委托 `analytics` | 原生批量 + 现有聚合 SQL |
|
||||
| 主库实现 | `logstore/postgres_store.go` | PG(按月分区)与 SQLite(普通表)共用 GORM |
|
||||
| Model | `internal/model/analytics` | 实体与批量 SQL,无 IO |
|
||||
| 入队 | `chwriter` / `risk_control` + `batchwriter` | flush 调 logstore `BatchInsert*`;CH 入队经 hooks |
|
||||
| 切换 | `of_log_db_switch` | 冻结 → `chwriter.Drain` → 逐表复制 → 翻转 |
|
||||
| 约束 | `logstore/imports_test.go` | apps 禁止 import `repository/analytics` |
|
||||
|
||||
`log_database` 只能是「随主库」或 `clickhouse`。`log_database` / `log_db_migration` 受保护。
|
||||
|
||||
## 新增一张日志表
|
||||
|
||||
1. **Model**(`internal/model/analytics`):`TableName` + `InsertColumns` / `BatchInsertSQL`。
|
||||
2. **三套 DDL**:CH `MergeTree` + `toYYYYMM`;PG `PARTITION BY RANGE(时间列)`(主键含分区键);SQLite 普通表。不要在主库建 CH 物化视图,聚合实时算。
|
||||
3. **挂到已有域或新接口**:能进 `AccessLogStore` / `ObservabilityStore` / `UserAccessLogStore` 就不要再拆包。新域才新增接口并放进 `Store`。
|
||||
4. **方法最少集**:`BatchInsert`(含 `ensureWritable`)、业务查询、`ListForMigration`、`MigrationRange`、`DeleteAll`、`DeleteBefore`、`EnsurePartitions`(仅 PG 预建)。
|
||||
5. **双实现**:CH 委托 `analyticsrepo`;GORM 共用一套,方言 SQL 放 `dialect_*.go`。零值 id 用 `idgen.NextUint64ID()`。
|
||||
6. **`buildStore`**:CH / GORM 两分支都挂上。
|
||||
7. **写入**:独立 `batchwriter`;`FlushFunc` → `logstore.Active`。节点日志/可观测走 `SetAccessLogHooks` / `SetObservabilityHooks`,不要让 apps 碰 `ChConn`。
|
||||
8. **切换任务**:`clearTarget` + `copy*` 增加该表;源数据不删,失败不翻转。
|
||||
9. **清理**:访问类走 `log_retention_days_*`;性能指标走 `metric_retention_days`。不要擅自共用错误的 TTL。
|
||||
10. **import-lint**:apps 新增对 `analytics` 或 `infra/persistence`(`batchwriter`/`idgen` 除外)的 import 必须失败。
|
||||
|
||||
## 禁止
|
||||
|
||||
- apps 直连 `analyticsrepo` / `db.ChConn` / `db.ChDB` 做日志读写
|
||||
- 只建 CH、不建主库回落
|
||||
- Handler 内逐条 `PrepareBatch`
|
||||
- 业务表塞进 logstore
|
||||
- 管理端改 `log_database` / `log_db_migration`
|
||||
|
||||
## 验证
|
||||
|
||||
```bash
|
||||
go test ./internal/repository/logstore ./internal/repository/analytics
|
||||
go test ./internal/apps/openflare/... ./internal/apps/admin/logs ./internal/apps/admin/status
|
||||
make swagger
|
||||
make code-check
|
||||
```
|
||||
|
||||
对照:`of_node_access_logs` 或 `w_user_access_logs` 的 model、三库 goose、`logstore` 双实现、`chwriter`/`risk_control` flush、`LogDBSwitchHandler`。
|
||||
Reference in New Issue
Block a user