mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 13:46:38 +08:00
Compare commits
26 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| fb08002e99 | |||
| f650214bbb | |||
| ba1c9222c2 | |||
| 076bf8b95c | |||
| 835c50dbaa | |||
| 42f7f47716 | |||
| 7d47db1f34 | |||
| 68d8f786cc | |||
| 9d93dc0b9f | |||
| 1a4a03a20d | |||
| 07e835c543 | |||
| 1f5bebd18a | |||
| fd62570431 | |||
| 484b49d79d | |||
| cffa009b8c | |||
| 4eced2b721 | |||
| ea7658815a | |||
| 3edcdb9e9f | |||
| 99f0f63b99 | |||
| 21fb303ef2 | |||
| 3aa4d98cd6 | |||
| 1f71c9f25b | |||
| 943818f7d4 | |||
| 23a5488203 | |||
| d99c5b7c43 | |||
| 33a1c32cf8 |
@@ -125,7 +125,7 @@ func startThingCacheInvalidationListener() {
|
||||
| Auth Source | `repository/auth_source_cache.go` | Otter | Redis JSON | `oauth:auth_source_invalidation` ✅ |
|
||||
| OAuth 用户/Token | `apps/oauth/cache.go` | 自研 map | Redis JSON | ❌ 无 pub/sub(历史债) |
|
||||
| 推送渠道 | `repository/push_channel.go` | 无 | Redis JSON | ❌ 仅 Redis Del |
|
||||
| Storage 驱动 | `internal/storage/storage.go` | RWMutex 快照 | — | `storage:config_invalidation` ✅ |
|
||||
| Storage 驱动 | `internal/infra/objectstore/storage.go` | RWMutex 快照 | — | `storage:config_invalidation` ✅ |
|
||||
|
||||
## 新增缓存工作流
|
||||
|
||||
@@ -210,7 +210,7 @@ make code-check
|
||||
## 相关文件
|
||||
|
||||
- L1 引擎:`pkg/cache/ram/cache.go`
|
||||
- DB/Redis 助手:`internal/db/redis.go`(`GetJSON`, `SetJSON`, `HGetJSON`, `PrefixedKey`)
|
||||
- DB/Redis 助手:`internal/infra/persistence/redis.go`(`GetJSON`, `SetJSON`, `HGetJSON`, `PrefixedKey`)
|
||||
- 金标准:`internal/repository/system_config_cache.go`
|
||||
- 上传元数据:`internal/apps/upload/cache/meta_cache.go`
|
||||
- Auth Source:`internal/repository/auth_source_cache.go`
|
||||
+14
-14
@@ -1,6 +1,6 @@
|
||||
---
|
||||
name: "clickhouse-batchwriter"
|
||||
description: "Wavelet 项目专用:当新增或修改 ClickHouse 批量写入、接入 internal/db/batchwriter、将业务域异步 flush 到分析表、迁移 risk_control/节点访问日志/可观测时序写入、或评估 async_insert 与背压策略时必须使用。本技能指导分层职责、各域独立 Writer 实例、repository 批量 API 与禁止写法。"
|
||||
description: "Wavelet 项目专用:当新增或修改 ClickHouse 批量写入、接入 internal/infra/persistence/batchwriter、将业务域异步 flush 到分析表、迁移 risk_control/节点访问日志/可观测时序写入、或评估 async_insert 与背压策略时必须使用。本技能指导分层职责、各域独立 Writer 实例、repository 批量 API 与禁止写法。"
|
||||
---
|
||||
|
||||
# ClickHouse 批量写入开发
|
||||
@@ -13,13 +13,13 @@ DDL 与表结构变更见 `database-migration` 技能;本技能只覆盖**运
|
||||
|
||||
| 层级 | 路径 | 职责 |
|
||||
| :--- | :--- | :--- |
|
||||
| 连接 | `internal/db/clickhouse.go` | `ChConn`(原生批量写)、`ChDB`(GORM 查询);禁止在业务包直接 `clickhouse.Open` |
|
||||
| 批量框架 | `internal/db/batchwriter/` | 泛型队列 + 按条数/时间 flush + 非阻塞入队 + 优雅停机;**各业务域独立实例** |
|
||||
| 连接 | `internal/infra/persistence/clickhouse.go` | `ChConn`(原生批量写)、`ChDB`(GORM 查询);禁止在业务包直接 `clickhouse.Open` |
|
||||
| 批量框架 | `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` |
|
||||
| 装配 | `internal/bootstrap/bootstrap.go` | 进程启动时调用 `Writer.Start`;初始化时需调用 `lifecycle.OnShutdown` 挂载停机钩子 |
|
||||
| 生命周期 | `internal/lifecycle/lifecycle.go` | 统一协调全局并发优雅停机,业务包无需在 `bootstrap.go` 中硬编码 `Stop` 逻辑 |
|
||||
| 装配 | `internal/platform/bootstrap/bootstrap.go` | 进程启动时调用 `Writer.Start`;初始化时需调用 `lifecycle.OnShutdown` 挂载停机钩子 |
|
||||
| 生命周期 | `internal/platform/lifecycle/lifecycle.go` | 统一协调全局并发优雅停机,业务包无需在 `bootstrap.go` 中硬编码 `Stop` 逻辑 |
|
||||
|
||||
**禁止**在 Handler / middleware 内直接 `db.ChConn.PrepareBatch`;**禁止**在 repository 内启动 goroutine 或维护全局 channel(队列生命周期由 apps + bootstrap 或专用 writer 包负责)。
|
||||
|
||||
@@ -68,7 +68,7 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush
|
||||
## 新增 ClickHouse 写入工作流
|
||||
|
||||
1. **Model**:在 `internal/model/analytics/` 定义 struct 与 `BatchInsertSQL()`(列顺序与 goose DDL 一致)。
|
||||
2. **Goose DDL**:在 `internal/db/migrator/goose/clickhouse/` 新增迁移(见 `database-migration`)。
|
||||
2. **Goose DDL**:在 `internal/infra/persistence/migrator/goose/clickhouse/` 新增迁移(见 `database-migration`)。
|
||||
3. **Repository**:实现 `BatchInsertX(ctx, []analyticsmodel.X) error`:
|
||||
- `len(items)==0` 直接返回
|
||||
- `db.ChConn == nil` 返回明确错误
|
||||
@@ -78,7 +78,7 @@ writer.Stop(stopCtx) // close 队列 + drain + 最终 flush
|
||||
- 业务路径 `TryEnqueue`;HTTP 背压用 `IsFull()`
|
||||
5. **测试**:
|
||||
- repository:mock `ChConn` 验证 `BatchInsertSQL` 与 append 列数
|
||||
- batchwriter:`go test ./internal/db/batchwriter`
|
||||
- batchwriter:`go test ./internal/infra/persistence/batchwriter`
|
||||
6. 运行 `make code-check`;有 API 变更时 `make swagger`。
|
||||
|
||||
## 背压与丢弃策略
|
||||
@@ -110,7 +110,7 @@ var globalChan chan any
|
||||
|
||||
## async_insert(补充,非主方案)
|
||||
|
||||
可在 `internal/db/clickhouse.go` 的 `Settings` 增加服务端异步写入作为第二层防护:
|
||||
可在 `internal/infra/persistence/clickhouse.go` 的 `Settings` 增加服务端异步写入作为第二层防护:
|
||||
|
||||
```go
|
||||
"async_insert": 1,
|
||||
@@ -122,7 +122,7 @@ var globalChan chan any
|
||||
## Bootstrap 装配示例
|
||||
|
||||
```go
|
||||
// internal/bootstrap/bootstrap.go(示意)
|
||||
// internal/platform/bootstrap/bootstrap.go(示意)
|
||||
var userAccessLogWriter *batchwriter.Writer[*analytics.UserAccessLog]
|
||||
|
||||
func RegisterAPI(ctx context.Context) {
|
||||
@@ -141,7 +141,7 @@ func RegisterAPI(ctx context.Context) {
|
||||
## 验证清单
|
||||
|
||||
```bash
|
||||
go test ./internal/db/batchwriter
|
||||
go test ./internal/infra/persistence/batchwriter
|
||||
go test ./internal/repository/analytics
|
||||
make code-check
|
||||
```
|
||||
@@ -153,11 +153,11 @@ make code-check
|
||||
|
||||
## 相关文件速查
|
||||
|
||||
- 框架:`internal/db/batchwriter/{config,writer,errs}.go`
|
||||
- 连接:`internal/db/clickhouse.go`
|
||||
- 框架:`internal/infra/persistence/batchwriter/{config,writer,errs}.go`
|
||||
- 连接:`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/lifecycle/lifecycle.go`
|
||||
- Bootstrap:`internal/bootstrap/bootstrap.go`
|
||||
- 生命周期管理器:`internal/platform/lifecycle/lifecycle.go`
|
||||
- Bootstrap:`internal/platform/bootstrap/bootstrap.go`
|
||||
+10
-10
@@ -1,17 +1,17 @@
|
||||
---
|
||||
name: "database-migration"
|
||||
description: "Wavelet 项目专用:当新增或修改数据库表结构、索引、初始化数据、系统配置 seed、模板 seed、默认管理员、goose SQL 迁移、internal/db/migrator、ClickHouse 分析库 DDL 或数据库升级流程时必须使用。本技能指导在 internal/db/migrator/goose 下编写 PostgreSQL/SQLite 双方言 SQL 迁移,以及在 goose/clickhouse 下编写 ClickHouse 单方言分析表迁移,并完成验证。"
|
||||
description: "Wavelet 项目专用:当新增或修改数据库表结构、索引、初始化数据、系统配置 seed、模板 seed、默认管理员、goose SQL 迁移、internal/infra/persistence/migrator、ClickHouse 分析库 DDL 或数据库升级流程时必须使用。本技能指导在 internal/infra/persistence/migrator/goose 下编写 PostgreSQL/SQLite 双方言 SQL 迁移,以及在 goose/clickhouse 下编写 ClickHouse 单方言分析表迁移,并完成验证。"
|
||||
---
|
||||
|
||||
# Wavelet 数据库升级操作指南
|
||||
|
||||
Wavelet 使用 `github.com/pressly/goose/v3` 执行 SQL 迁移。迁移入口是 `internal/db/migrator.Migrate()`,SQL 文件嵌入在二进制中。
|
||||
Wavelet 使用 `github.com/pressly/goose/v3` 执行 SQL 迁移。迁移入口是 `internal/infra/persistence/migrator.Migrate()`,SQL 文件嵌入在二进制中。
|
||||
|
||||
## 基本规则
|
||||
|
||||
- SQL 迁移文件放在:
|
||||
- `internal/db/migrator/goose/postgres/`
|
||||
- `internal/db/migrator/goose/sqlite/`
|
||||
- `internal/infra/persistence/migrator/goose/postgres/`
|
||||
- `internal/infra/persistence/migrator/goose/sqlite/`
|
||||
- PostgreSQL 和 SQLite 必须使用同一个版本号、同一个语义文件名。
|
||||
- 迁移文件使用 goose SQL 标记:
|
||||
|
||||
@@ -51,7 +51,7 @@ Wavelet 使用 `github.com/pressly/goose/v3` 执行 SQL 迁移。迁移入口是
|
||||
7. 至少运行:
|
||||
|
||||
```bash
|
||||
go test ./internal/db/migrator
|
||||
go test ./internal/infra/persistence/migrator
|
||||
go test ./internal/model ./internal/apps/config ./internal/apps/admin/system_config
|
||||
make code-check
|
||||
```
|
||||
@@ -89,10 +89,10 @@ ClickHouse 是**辅助 OLAP 存储**,与 PostgreSQL/SQLite 主库**完全独
|
||||
|
||||
| 路径 | 职责 |
|
||||
| :--- | :--- |
|
||||
| `internal/db/migrator/goose/clickhouse/` | **唯一** ClickHouse DDL 来源(goose SQL,嵌入二进制) |
|
||||
| `internal/infra/persistence/migrator/goose/clickhouse/` | **唯一** ClickHouse DDL 来源(goose SQL,嵌入二进制) |
|
||||
| `internal/model/analytics/` | 分析表 Go model,列名须与 goose DDL 一致 |
|
||||
| `internal/repository/analytics/` | 所有 ClickHouse 读写(批量写入、查询、聚合) |
|
||||
| `internal/db/clickhouse.go` | 连接初始化(`ChConn` 原生批量、`ChDB` GORM 查询) |
|
||||
| `internal/infra/persistence/clickhouse.go` | 连接初始化(`ChConn` 原生批量、`ChDB` GORM 查询) |
|
||||
|
||||
### 迁移入口与版本表
|
||||
|
||||
@@ -116,16 +116,16 @@ ClickHouse 是**辅助 OLAP 存储**,与 PostgreSQL/SQLite 主库**完全独
|
||||
按以下顺序落地,避免列名或类型漂移:
|
||||
|
||||
1. **Model**:在 `internal/model/analytics/` 定义 struct,`gorm:"column:..."` 与 DDL 列名一一对应;实现 `TableName()`,批量写入表可提供 `InsertColumns()` / `BatchInsertSQL()`。
|
||||
2. **Goose SQL**:在 `internal/db/migrator/goose/clickhouse/` 新增递增版本文件(格式同主库,如 `YYYYMMDDNNNN_create_xxx.sql`),编写 `-- +goose Up` / `-- +goose Down`。
|
||||
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/db/batchwriter` 各域独立实例异步 flush(详见 `clickhouse-batchwriter` 技能),`FlushFunc` 只调 repository `BatchInsert*`;管理端统计 API 只读 repository,不触达 DDL。
|
||||
4. **Apps**:在 `internal/apps/<domain>/` 编排采集与入队;高频写入通过 `internal/infra/persistence/batchwriter` 各域独立实例异步 flush(详见 `clickhouse-batchwriter` 技能),`FlushFunc` 只调 repository `BatchInsert*`;管理端统计 API 只读 repository,不触达 DDL。
|
||||
|
||||
### ClickHouse 验证
|
||||
|
||||
至少运行:
|
||||
|
||||
```bash
|
||||
go test ./internal/db/migrator
|
||||
go test ./internal/infra/persistence/migrator
|
||||
go test ./internal/repository/analytics
|
||||
make code-check
|
||||
```
|
||||
@@ -15,7 +15,7 @@ Wavelet 将「对象存储」与「上传业务」分为两层,**禁止混用
|
||||
|
||||
| 层级 | 包路径 | 职责 | 业务是否直接调用 |
|
||||
| :--- | :--- | :--- | :--- |
|
||||
| **对象存储引擎** | `internal/storage` | `Backend` 接口:`Put` / `Get` / `Delete` / `Test`;按配置切换 Local / S3 / R2 / OSS / WebDAV | **禁止**(仅 upload 域内部使用) |
|
||||
| **对象存储引擎** | `internal/infra/objectstore` | `Backend` 接口:`Put` / `Get` / `Delete` / `Test`;按配置切换 Local / S3 / R2 / OSS / WebDAV | **禁止**(仅 upload 域内部使用) |
|
||||
| **上传域服务** | `internal/apps/upload` | `w_uploads` 记录、权限、秒传、统计、文件服务、`upload.Ingest` | **必须** |
|
||||
| **上传 HTTP 入口** | `internal/apps/upload/handler` | `POST /api/v1/upload` 等 multipart 接口 | 前端 / 用户侧上传 |
|
||||
| **文件访问** | `internal/apps/upload/filesrv` | `GET /f/:id` 流式响应、访问控制、图片 WebP 压缩 | 展示 / 下载 |
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user