diff --git a/hzh/Gen2D/00-索引.md b/hzh/Gen2D/00-索引.md index d349dbb..6986ab4 100644 --- a/hzh/Gen2D/00-索引.md +++ b/hzh/Gen2D/00-索引.md @@ -1,62 +1,89 @@ +--- +tags: [index, architecture, gen2d, overview] +create time: 2026-06-03 09:00 +--- + # Gen2D 架构讲解 — 答辩文档索引 > AI 驱动的 2D 游戏素材生成工具 -## 文档导航 +## 概述 + +本文档是 Gen2D 系统架构讲解的导航索引,覆盖从整体分层架构到具体技术组件的 15 个核心知识点。通过本文档可以快速定位到任意模块的详细解读。 + +## 正文 + +### 文档导航 | # | 文件 | 主题 | 核心要点 | |---|------|------|----------| -| 1 | [01-system-overview](01-system-overview.md) | 系统总览 | Gin → Handler → Service → 基础设施 → 外部依赖 | -| 2 | [02-worker-pool](02-worker-pool.md) | 协程池 | 有界并发、per-user 限流、背压、优雅关闭 | -| 3 | [03-task-queue](03-task-queue.md) | 任务队列 | 可插拔接口 + Memory / RabbitMQ 双实现 | -| 4 | [04-rabbitmq](04-rabbitmq.md) | RabbitMQ 集成 | AMQP 连接、持久化消息、重试/死信流程 | -| 5 | [05-generation-pipeline](05-generation-pipeline.md) | 生成管线 | Eino 4 阶段图 + 质量回退 + 降级路径 | -| 6 | [06-sprite-processing](06-sprite-processing.md) | 精灵图处理 | 背景移除 → 投影切割 → 后处理 → GIF 预览 | -| 7 | [07-observability](07-observability.md) | 可观测性 | 35 Prometheus 指标 + 3 Grafana 仪表盘 + 10 告警 | -| 8 | [08-sse-push](08-sse-push.md) | SSE 实时推送 | EventBus → SSEHandler → 浏览器 EventSource | -| 9 | [09-rate-limiting](09-rate-limiting.md) | 限流 | Redis Lua 令牌桶 + 双层限流 + Fail-Open | -| 10 | [10-middleware-chain](10-middleware-chain.md) | 中间件链 | Logger → Recovery → Metrics → Auth → RateLimit | -| 11 | [11-consumer-producer](11-consumer-producer.md) | Consumer-Producer 桥接 | TaskQueue → Consumer → WorkerPool 解耦 | -| 12 | [12-three-tier-fallback](12-three-tier-fallback.md) | 三级降级策略 | Queue → Pool → Legacy 降级链 | -| 13 | [13-prompt-engineering](13-prompt-engineering.md) | 标签驱动提示词 | 40+ 标签映射 + 风格一致性 | -| 14 | [14-deployment](14-deployment.md) | 部署架构 | Docker Compose + Prometheus + Grafana | -| 15 | [15-config-cascade](15-config-cascade.md) | 配置级联 | Viper 三层配置:YAML → ENV → Default | +| 1 | [[01-系统总览]] | 系统总览 | Gin → Handler → Service → 基础设施 → 外部依赖 | +| 2 | [[02-协程池]] | 协程池 | 有界并发、per-user 限流、背压、优雅关闭 | +| 3 | [[03-任务队列]] | 任务队列 | 可插拔接口 + Memory / RabbitMQ 双实现 | +| 4 | [[04-RabbitMQ集成]] | RabbitMQ 集成 | AMQP 连接、持久化消息、重试/死信流程 | +| 5 | [[05-生成管线]] | 生成管线 | Eino 4 阶段图 + 质量回退 + 降级路径 | +| 6 | [[06-精灵图处理]] | 精灵图处理 | 背景移除 → 投影切割 → 后处理 → GIF 预览 | +| 7 | [[07-可观测性]] | 可观测性 | 35 Prometheus 指标 + 3 Grafana 仪表盘 + 10 告警 | +| 8 | [[08-SSE实时推送]] | SSE 实时推送 | EventBus → SSEHandler → 浏览器 EventSource | +| 9 | [[09-限流]] | 限流 | Redis Lua 令牌桶 + 双层限流 + Fail-Open | +| 10 | [[10-中间件链]] | 中间件链 | Logger → Recovery → Metrics → Auth → RateLimit | +| 11 | [[11-Consumer-Producer桥接]] | Consumer-Producer 桥接 | TaskQueue → Consumer → WorkerPool 解耦 | +| 12 | [[12-三级降级策略]] | 三级降级策略 | Queue → Pool → Legacy 降级链 | +| 13 | [[13-标签驱动提示词工程]] | 标签驱动提示词 | 40+ 标签映射 + 风格一致性 | +| 14 | [[14-部署架构]] | 部署架构 | Docker Compose + Prometheus + Grafana | +| 15 | [[15-配置级联机制]] | 配置级联 | Viper 三层配置:YAML → ENV → Default | -## 架构总览图 +> [!tip] 阅读顺序建议 +> +> 建议按编号顺序阅读:先理解 [[01-系统总览]] 的全局架构,再深入各个具体组件(协程池、任务队列、RabbitMQ),最后学习运行时相关话题(SSE 推送、可观测性、降级策略)。 -``` -┌─────────────────────────────────────────────────────────────────┐ -│ 浏览器 (React) │ -│ EventSource ← SSE ← EventBus ← Pipeline Progress │ -└────────────────────────────┬────────────────────────────────────┘ - │ HTTP -┌────────────────────────────▼────────────────────────────────────┐ -│ Gin HTTP Server │ -│ Logger → Recovery → Metrics → Auth → RateLimit → Handler │ -└────────────────────────────┬────────────────────────────────────┘ - │ -┌────────────────────────────▼────────────────────────────────────┐ -│ Handler 层 │ -│ GenerateHandler / SSEHandler / PromptHandler │ -└──────┬─────────────────────┬─────────────────────┬──────────────┘ - │ │ │ -┌──────▼──────┐ ┌─────────▼─────────┐ ┌─────▼──────┐ -│ TaskQueue │ │ WorkerPool │ │ EventBus │ -│ Memory/RMQ │───→│ NumCPU*4 并发 │ │ Pub/Sub │ -└─────────────┘ │ Per-user 限流 │ └────────────┘ - └─────────┬─────────┘ - │ - ┌─────────▼─────────┐ - │ Eino Pipeline │ - │ 4 阶段生成管线 │ - └─────────┬─────────┘ - │ - ┌──────────────┼──────────────┐ - ▼ ▼ ▼ - Image API Quality Check SplitSprite - (DALL-E 3) (GPT-4o) + GIF Maker +### 架构总览图 + +```mermaid +graph TB + subgraph Browser["浏览器 (React)"] + UI["三栏工作台"] + SE["EventSource ← SSE"] + end + + subgraph Gateway["Gin HTTP Server"] + MW["中间件链
Logger / Recovery / Metrics / Auth / RateLimit"] + end + + subgraph Handler["Handler 层"] + GH["GenerateHandler"] + SH["SSEHandler"] + PH["PromptHandler"] + end + + subgraph Infra["基础设施层"] + TQ["TaskQueue
Memory / RabbitMQ"] + WP["WorkerPool
NumCPU*4 并发
Per-user 限流"] + EB["EventBus
Pub/Sub"] + end + + subgraph Pipeline["Service 层"] + EP["Eino Pipeline
4 阶段生成管线"] + end + + subgraph External["外部依赖"] + IMG["Image API
DALL-E 3"] + QC["Quality Check
GPT-4o"] + SP["SplitSprite + GIF Maker"] + end + + UI -->|"HTTP"| GW + SE --> UI + GW --> MW + MW --> Handler + GH --> TQ + GH --> WP + SH --> EB + TQ --> WP + WP --> EP + EP --> IMG & QC & SP ``` -## 配套图表 +### 配套图表 每份文档开头使用 Mermaid 流程图,可直接在 Markdown 渲染器中查看。 diff --git a/hzh/Gen2D/01-系统总览.md b/hzh/Gen2D/01-系统总览.md index 9bfaffb..9c5464f 100644 --- a/hzh/Gen2D/01-系统总览.md +++ b/hzh/Gen2D/01-系统总览.md @@ -1,8 +1,19 @@ -# 01 - 系统总览 +--- +tags: [architecture, system-design, go, gin, dependency-injection, viper, config] +create time: 2026-06-03 10:00 +--- -> **一句话概括**:Gen2D 采用经典分层架构,Gin HTTP Server -> Handler -> Service -> 基础设施 -> 外部依赖,各层职责清晰、可独立替换。 +# 01. 系统总览 -## 架构全景 +## 概述 + +Gen2D 采用经典分层架构,数据流自上而下贯穿 Gin HTTP Server → Handler → Service → 基础设施 → 外部依赖,各层职责清晰、可独立替换。 + +> **一句话概括**:分层架构,职责分离,轻量 DI,零配置可启动。 + +## 正文 + +### 架构全景 ```mermaid graph TB @@ -57,7 +68,7 @@ graph TB Pipeline -->|"Progress"| EB ``` -## 分层详解 +### 分层详解 | 层级 | 目录 | 核心职责 | 代表组件 | |------|------|----------|----------| @@ -70,9 +81,9 @@ graph TB | **精灵处理** | `pkg/splitsprite/`, `pkg/gifmaker/` | 精灵表切割、GIF 预览生成 | `splitsprite.Process`, `gifmaker.Encode` | | **配置** | `internal/config/` | YAML 加载、环境变量绑定、默认值 | `config.Load()` | -## 依赖注入模式 +### 依赖注入模式 -Gen2D 采用轻量级的 **Set\*/Init\* 函数注入** 模式,避免引入 DI 框架。 +Gen2D 采用轻量级的 **Set*/Init* 函数注入** 模式,避免引入 DI 框架。 ``` main.go 中的注入链路: @@ -93,9 +104,11 @@ eventbus.Init() → 初始化事件总线 - `service.InitImageGenConfig()` 将配置缓存为包级变量,避免在函数签名中传递大量参数 - 每个 `Set*` 函数对应一个包级全局变量,简单但足够清晰 -> :bulb: **为什么不用 Wire / Fx?** 项目规模可控,`cmd/main.go` 约 220 行即可完成全部注入,框架级 DI 的复杂度收益比不高。 +> [!tip] 为什么不用 Wire / Fx? +> +> 项目规模可控,`cmd/main.go` 约 220 行即可完成全部注入,框架级 DI 的复杂度收益比不高。 -## 配置级联 +### 配置级联 Gen2D 使用 Viper 实现三层配置覆盖,优先级从高到低: @@ -130,40 +143,50 @@ graph LR | `workerpool` | `GEN2D_WORKERPOOL_*` | workers, queue_size, max_per_user | | `taskqueue` | `GEN2D_TASKQUEUE_*` | driver, memory.buffer_size, rabbitmq.* | -## 请求全链路 +### 请求全链路 一个素材生成请求的完整生命周期: -``` -1. 浏览器 → POST /api/v1/generate -2. Gin 中间件链 → Logger → Recovery → Metrics → Auth → RateLimit -3. Handler.Generate() - ├── 参数绑定 & 校验 - ├── 保存任务到 MySQL(status=pending) - ├── 构建 TaskMessage - └── 提交到 TaskQueue(优先)或 WorkerPool(fallback) -4. 返回 { taskId } 给浏览器 -5. Consumer 从 TaskQueue 消费消息 -6. Consumer → WorkerPool.Submit() -7. WorkerPool 的 Worker 执行任务: - ├── 注入 ProgressReporter 到 context - ├── RunPipeline() — Eino Graph 执行 - │ ├── PromptOptimizer(合并风格 + 技术参数) - │ ├── AssetGenerator(调用 AI 出图 API) - │ ├── QualitySupervisor(质检,最多重试 3 次) - │ └── FormatAdapter(精灵表切割 + GIF 预览) - ├── 上传素材到七牛云 - └── 更新 MySQL 任务状态 -8. EventBus.Publish() → SSE 推送给浏览器 -9. 浏览器通过 EventSource 实时接收进度 +```mermaid +graph LR + BROWSER["浏览器
POST /api/v1/generate"] --> GIN["Gin中间件链"] + GIN --> HANDLER["Handler.Generate"] + HANDLER --> VALIDATE["参数绑定校验"] + HANDLER --> MYSQL["保存任务MySQL
status=pending"] + HANDLER --> MSG["构建TaskMessage"] + MSG --> TQ{"提交目标"} + TQ -->|"优先"| TASKQUEUE["TaskQueue队列排队"] + TQ -->|"Fallback"| WORKERPOOL["WorkerPool直连"] + VALIDATE --> TQ + MYSQL --> TQ + TASKQUEUE --> CONSUMER["Consumer消费"] + WORKERPOOL --> SUBMIT["WorkerPool.Submit"] + CONSUMER --> SUBMIT + SUBMIT --> WORKER["Worker执行"] + WORKER --> PROGRESS["注入ProgressReporter"] + WORKER --> PIPELINE["Eino Pipeline执行"] + PIPELINE --> PROMPT["PromptOptimizer"] + PIPELINE --> ASSET["AssetGenerator"] + PIPELINE --> QUALITY["QualitySupervisor
质检最多重试3次"] + PIPELINE --> FORMAT["FormatAdapter
精灵表切割GIF预览"] + PROGRESS --> UPLOAD["上传素材七牛云"] + QUALITY --> UPLOAD + UPLOAD --> STATUS["更新MySQL状态"] + STATUS --> EVENTBUS["EventBus.Publish"] + EVENTBUS --> SSE["SSE推送"] + SSE --> BROWSER_SSE["浏览器EventSource
实时接收进度"] ``` -## 关键设计决策 +### 关键设计决策 + +> [!question] 思考:为什么 Gen2D 选择了异步任务 + SSE 推送的组合? +> +> 如果直接同步调用 Eino Pipeline,一个生成请求可能要等待 10~120 秒。在 HTTP 模型下,长时间占用的连接会耗尽服务器的并发能力。**异步提交 + SSE 推送**把「等待时间」从连接持有中解放出来——客户端收到 taskId 后可以自由离开,后续通过 SSE 长连接接收进度更新。这也是 Web 应用在 AI 场景下的标准模式。 | 决策 | 选择 | 理由 | |------|------|------| | 分层架构 | Handler-Service-Infra 三层 | 解耦各层职责,便于独立测试和替换 | -| 依赖注入 | Set\*/Init\* 函数 | 轻量、零依赖,项目规模可控 | +| 依赖注入 | Set*/Init* 函数 | 轻量、零依赖,项目规模可控 | | 配置管理 | Viper 三层级联 | 容器化友好,零配置可启动 | | 异步任务 | 提交-队列-消费-执行 | API 快速返回,长任务不阻塞请求 | | 事件推送 | EventBus + SSE | 比 WebSocket 轻量,HTTP 原生支持 | @@ -171,8 +194,8 @@ graph LR ## 关联文档 -- [索引](00-index.md) — 文档导航与架构总览图 -- [协程池](02-worker-pool.md) — 有界并发与 per-user 限流 -- [任务队列](03-task-queue.md) — 可插拔队列接口与双实现 -- [RabbitMQ 集成](04-rabbitmq.md) — 持久化消息与重试机制 -- [生成管线](05-generation-pipeline.md) — Eino 4 阶段管线与质量回退 +- [[00-索引]] — 文档导航与架构总览图 +- [[02-协程池]] — 有界并发与 per-user 限流 +- [[03-任务队列]] — 可插拔队列接口与双实现 +- [[04-RabbitMQ集成]] — 持久化消息与重试机制 +- [[05-生成管线]] — Eino 4 阶段管线与质量回退 diff --git a/hzh/Gen2D/02-协程池.md b/hzh/Gen2D/02-协程池.md index 9d6af57..2dcab18 100644 --- a/hzh/Gen2D/02-协程池.md +++ b/hzh/Gen2D/02-协程池.md @@ -1,22 +1,33 @@ -# 02 - 协程池 (Worker Pool) +--- +tags: [concurrency, goroutine-pool, backpressure, go, worker-pattern, rate-limiting] +create time: 2026-06-03 10:05 +--- + +# 02. 协程池 (Worker Pool) + +## 概述 + +有界并发协程池,以 `NumCPU*4` 个 Worker 并行处理任务,配合 per-user 限流和背压保护,支持优雅关闭。 > **一句话概括**:有界并发协程池,`NumCPU*4` workers,per-user 限流,背压保护,优雅关闭。 -## 工作流 +## 正文 + +### 工作流 ```mermaid graph TB - subgraph Submit["Submit() 入口"] + subgraph Submit["Submit入口"] CHECK_CLOSED{"池已关闭?"} - CHECK_USER{"per-user 限流
active >= max?"} - CHECK_FULL{"channel 满?"} + CHECK_USER{"per-user限流
active >= max?"} + CHECK_FULL{"channel满?"} end subgraph Channel["有界任务队列"] TASK_CHAN["chan Task
capacity = queueSize"] end - subgraph Workers["Worker 协程"] + subgraph Workers["Worker协程"] W1["Worker 0"] W2["Worker 1"] W3["Worker ..."] @@ -25,7 +36,7 @@ graph TB subgraph Execute["任务执行"] TIMEOUT["context.WithTimeout
10 min"] - FN["task.Fn(ctx)"] + FN["task.Fnctx"] RELEASE["释放用户槽位"] end @@ -33,13 +44,13 @@ graph TB CHECK_CLOSED -->|否| CHECK_USER CHECK_USER -->|"ErrUserLimitReached"| REJECT CHECK_USER -->|通过| CHECK_FULL - CHECK_FULL -->|"ErrPoolFull (HTTP 503)"| REJECT + CHECK_FULL -->|"ErrPoolFull HTTP 503"| REJECT CHECK_FULL -->|通过| TASK_CHAN TASK_CHAN --> W1 & W2 & W3 & WN W1 & W2 & W3 & WN --> TIMEOUT --> FN --> RELEASE ``` -## Pool 结构体 +### Pool 结构体 ```go type Pool struct { @@ -70,7 +81,7 @@ type Pool struct { | `taskQueue` | `chan Task` | buffered channel | 有界队列,固定容量 | | `userActive` | `map[string]int` | — | 记录每用户活跃任务数 | -## Functional Options 模式 +### Functional Options 模式 协程池采用 Functional Options 模式进行配置,开箱即用、可选覆盖: @@ -88,9 +99,11 @@ pool := workerpool.New( | `WithQueueSize(n)` | `100` | `n < 1` 时强制为 1 | 有界缓冲,满时触发背压 | | `WithMaxPerUser(n)` | `2` | `n < 1` 时置 0(不限制) | 防止单用户占满池 | -> :bulb: **为什么选择 Functional Options?** 可选参数天然为零值时保持默认,新增配置项无需修改 `New()` 签名,调用方按需指定。 +> [!tip] 为什么选择 Functional Options? +> +> 可选参数天然为零值时保持默认,新增配置项无需修改 `New()` 签名,调用方按需指定。 -## 有界并发 +### 有界并发 协程池的核心是 `make(chan Task, queueSize)` 创建的 **有界缓冲 channel**。 @@ -112,7 +125,7 @@ default: // 队列满,背压 } ``` -## Per-user 限流 +### Per-user 限流 每个用户同时执行的任务数受到 `maxPerUser` 限制,防止单用户占满整个池。 @@ -134,55 +147,50 @@ Submit() 调用流程: | Execute | — | Worker 从 channel 取出后开始执行 | | Complete | `defer userActive[userID]--` | 任务完成或失败时释放 | -> :warning: **槽位预留时机**:在 `Submit()` 而非 `worker()` 中预留,确保 channel 满时不会出现"槽位已分配但任务未入队"的不一致状态。 +> [!warning] 槽位预留时机 +> +> 在 `Submit()` 而非 `worker()` 中预留,确保 channel 满时不会出现"槽位已分配但任务未入队"的不一致状态。 -## 背压保护 +### 背压保护 当任务队列已满时,协程池通过 `select + default` 实现非阻塞拒绝: -``` -队列满(channel 已达 capacity) - │ - ▼ -select 进入 default 分支 - │ - ├── 回滚用户槽位(如果有) - ├── 原子递增 RejectedTasks - └── 返回 ErrPoolFull - │ - ▼ - Handler 层映射为 HTTP 503 Service Unavailable - 响应体:"系统繁忙,请稍后重试" +```mermaid +graph TB + FULL["队列满channel已达capacity"] --> SELECT["select进入default分支"] + SELECT --> ROLLBACK["回滚用户槽位如果有"] + SELECT --> INC["原子递增RejectedTasks"] + SELECT --> RETURN["返回ErrPoolFull"] + RETURN --> HTTP503["Handler层映射为HTTP 503
响应体系统繁忙请稍后重试"] ``` **背压 vs 阻塞**: | 策略 | 行为 | 适用场景 | |------|------|----------| -| 非阻塞拒绝(Gen2D) | 立即返回错误 | 用户交互型 API,快速失败 | +| 非阻塞拒绝Gen2D | 立即返回错误 | 用户交互型 API,快速失败 | | 阻塞等待 | 阻塞直到有空位 | 批处理系统,不能丢任务 | Gen2D 选择非阻塞拒绝,因为用户期望快速得到反馈,而非无限等待。 -## 优雅关闭 +> [!tip] 设计权衡 +> +> 选择「非阻塞拒绝」意味着可能丢失用户的提交意图。但在 AI 生成场景中,用户可以 **重新点击提交按钮**——这种短暂的操作成本远低于服务器因连接堆积导致的雪崩风险。这就是经典的 **「可用性 > 一致性」** 抉择。 + +### 优雅关闭 协程池支持优雅关闭,确保正在执行的任务有时间完成: -``` -收到 SIGINT / SIGTERM - │ - ▼ -consumeCancel() ← 停止消费者,不再接收新任务 - │ - ▼ -pool.Shutdown(ctx) ← 30 秒超时 - │ - ├── closed.Swap(true) ← 停止接收新任务 - ├── cancel() ← 通知 Worker 停止取任务 - ├── wg.Wait() ← 等待所有 Worker 退出 - │ - ├── 成功 → 日志 "workerpool shutdown gracefully" - └── 超时 → 日志 "workerpool shutdown timeout" +```mermaid +graph TB + SIG["收到SIGINT / SIGTERM"] --> STOP_CONSUME["consumeCancel
停止消费者不再接收新任务"] + STOP_CONSUME --> SHUTDOWN["pool.Shutdownctx
30秒超时"] + SHUTDOWN --> SWAP["closed.Swaptrue
停止接收新任务"] + SHUTDOWN --> CANCEL["cancel
通知Worker停止取任务"] + SHUTDOWN --> WAIT["wg.Wait
等待所有Worker退出"] + WAIT --> SUCCESS{"成功?"} + SUCCESS -->|是| LOG_OK["日志workerpool shutdown gracefully"] + SUCCESS -->|否| LOG_TIMEOUT["日志workerpool shutdown timeout"] ``` **Shutdown 返回值**: @@ -192,7 +200,7 @@ pool.Shutdown(ctx) ← 30 秒超时 | `true` | 所有任务正常完成 | | `false` | 超时,部分任务可能丢失 | -## Metrics 指标 +### Metrics 指标 协程池内置原子计数器,支持运行时观测: @@ -207,7 +215,7 @@ pool.Shutdown(ctx) ← 30 秒超时 所有指标通过 `Metrics()` 方法返回只读快照,同时上报 Prometheus。 -## Task 结构体 +### Task 结构体 ```go type Task struct { @@ -220,19 +228,13 @@ type Task struct { 每个任务携带超时 context(默认 10 分钟),从池的根 context 派生,确保优雅关闭时能取消正在执行的任务。 -## 与 TaskQueue 的协作 +### 与 TaskQueue 的协作 -``` -TaskQueue(全局排队) - │ - ▼ -Consumer(消费消息) - │ - ▼ -WorkerPool.Submit()(单机并发控制) - │ - ▼ -Worker 执行 Pipeline +```mermaid +graph LR + TQ["TaskQueue全局排队"] --> CONSUMER["Consumer消费消息"] + CONSUMER --> WP["WorkerPoolSubmit单机并发控制"] + WP --> WORKER["Worker执行Pipeline"] ``` | 组件 | 职责 | 范围 | @@ -243,7 +245,7 @@ Worker 执行 Pipeline ## 关联文档 -- [索引](00-index.md) — 文档导航与架构总览图 -- [系统总览](01-system-overview.md) — 分层架构与依赖注入 -- [任务队列](03-task-queue.md) — 可插拔队列接口 -- [Consumer-Producer 桥接](00-index.md) — TaskQueue 到 WorkerPool 的解耦 +- [[00-索引]] — 文档导航与架构总览图 +- [[01-系统总览]] — 分层架构与依赖注入 +- [[03-任务队列]] — 可插拔队列接口 +- [[11-Consumer-Producer桥接]] — TaskQueue 到 WorkerPool 的解耦 diff --git a/hzh/Gen2D/03-任务队列.md b/hzh/Gen2D/03-任务队列.md index 2a706be..830ee85 100644 --- a/hzh/Gen2D/03-任务队列.md +++ b/hzh/Gen2D/03-任务队列.md @@ -1,19 +1,30 @@ -# 03 - 任务队列 (Task Queue) +--- +tags: [task-queue, plugin-architecture, memory-queue, rabbitmq, go, interface-pattern] +create time: 2026-06-03 10:10 +--- + +# 03. 任务队列 (Task Queue) + +## 概述 + +可插拔任务队列接口,支持 Memory 和 RabbitMQ 双实现,通过工厂模式一行切换,满足不同部署环境的需求。 > **一句话概括**:可插拔任务队列接口,支持 Memory 和 RabbitMQ 双实现,通过工厂模式一行切换。 -## 架构设计 +## 正文 + +### 架构设计 ```mermaid graph TB subgraph Producer["生产者"] - HANDLER["Handler.Generate()"] + HANDLER["Handler.Generate"] end - subgraph Interface["TaskQueue 接口"] - SUBMIT["Submit(ctx, msg)"] - CONSUME["Consume(ctx, handler)"] - CLOSE["Close()"] + subgraph Interface["TaskQueue接口"] + SUBMIT["Submitctx msg"] + CONSUME["Consumectx handler"] + CLOSE["Close"] end subgraph Memory["MemoryQueue"] @@ -36,7 +47,7 @@ graph TB RabbitMQ --- RMQ_PUB ``` -## TaskQueue 接口 +### TaskQueue 接口 ```go type TaskQueue interface { @@ -52,7 +63,7 @@ type TaskQueue interface { | `Consume` | 持续消费任务,直到 ctx 取消 | handler 返回错误时,内存队列丢弃,RabbitMQ NACK 重试 | | `Close` | 关闭连接,释放资源 | 返回 `errors.Join` 聚合错误 | -## 工厂模式 +### 工厂模式 通过配置驱动,一行切换队列实现: @@ -76,7 +87,7 @@ func New(cfg config.TaskQueueConfig) (TaskQueue, error) { | `"memory"` (默认) | `MemoryQueue` | 单机开发、演示环境 | | `"rabbitmq"` | `RabbitMQQueue` | 多机生产部署 | -## TaskMessage 消息结构 +### TaskMessage 消息结构 ```go type TaskMessage struct { @@ -95,7 +106,7 @@ type TaskMessage struct { - `RetryCount` 供 RabbitMQ 实现判断是否超过最大重试次数 - JSON 序列化,兼容内存队列和 RabbitMQ 两种传输 -## MemoryQueue 实现 +### MemoryQueue 实现 基于 Go channel 的内存队列,零外部依赖。 @@ -109,7 +120,7 @@ type MemoryQueue struct { } ``` -### Submit +#### Submit ```go func (q *MemoryQueue) Submit(ctx context.Context, msg TaskMessage) error { @@ -124,9 +135,11 @@ func (q *MemoryQueue) Submit(ctx context.Context, msg TaskMessage) error { } ``` -> :bulb: **阻塞语义**:MemoryQueue 的 Submit 是阻塞的——当 buffer 满时,调用方会阻塞直到有空位或 ctx 取消。这与 WorkerPool 的非阻塞拒绝形成对比。 +> [!tip] 阻塞语义 +> +> MemoryQueue 的 Submit 是阻塞的——当 buffer 满时,调用方会阻塞直到有空位或 ctx 取消。这与 WorkerPool 的非阻塞拒绝形成对比。 -### Consume +#### Consume ```go func (q *MemoryQueue) Consume(ctx context.Context, handler func(TaskMessage) error) error { @@ -155,11 +168,13 @@ func (q *MemoryQueue) Consume(ctx context.Context, handler func(TaskMessage) err | 背压 | 阻塞直到有空位 | | 依赖 | 零外部依赖 | -> :warning: **可接受的任务丢失**:AI 生成任务可以重新提交,进程重启丢失排队中的任务是可接受的权衡。 +> [!warning] 可接受的任务丢失 +> +> AI 生成任务可以重新提交,进程重启丢失排队中的任务是可接受的权衡。但如果你的业务场景中 **任务不可重放**(比如支付指令),则必须选择 RabbitMQ 等持久化实现。 -## RabbitMQQueue 实现 +### RabbitMQQueue 实现 -基于 AMQP 的持久化消息队列,支持手动 ACK 和重试。详见 [04-rabbitmq](04-rabbitmq.md)。 +基于 AMQP 的持久化消息队列,支持手动 ACK 和重试。详见 [[04-RabbitMQ集成]]。 **RabbitMQQueue 特性**: @@ -170,7 +185,7 @@ func (q *MemoryQueue) Consume(ctx context.Context, handler func(TaskMessage) err | 背压 | 发布失败时返回错误(不阻塞) | | 依赖 | 需要 RabbitMQ 服务 | -## 双实现对比 +### 双实现对比 | 维度 | MemoryQueue | RabbitMQQueue | |------|-------------|---------------| @@ -182,7 +197,7 @@ func (q *MemoryQueue) Consume(ctx context.Context, handler func(TaskMessage) err | **适用** | 开发/演示 | 生产环境 | | **消息丢失** | 进程重启丢失 | 服务重启不丢失 | -## Consumer 桥接 +### Consumer 桥接 Consumer 从 TaskQueue 消费消息,提交到 WorkerPool 执行,实现队列与并发控制的解耦: @@ -201,12 +216,14 @@ func (c *Consumer) Start(ctx context.Context) error { } ``` -``` -TaskQueue → Consumer → WorkerPool → Pipeline - 全局排队 桥接 单机并发 业务逻辑 +```mermaid +graph LR + TQ["TaskQueue全局排队"] --> CONSUMER["Consumer桥接"] + CONSUMER --> WP["WorkerPool单机并发"] + WP --> PIPELINE["Pipeline业务逻辑"] ``` -## 三级降级策略 +### 三级降级策略 Gen2D 在 `cmd/main.go` 中实现了三级降级链: @@ -229,7 +246,7 @@ if taskQueue != nil { } ``` -## Metrics 指标 +### Metrics 指标 | Prometheus 指标 | 类型 | Label | 说明 | |-----------------|------|-------|------| @@ -241,7 +258,7 @@ if taskQueue != nil { ## 关联文档 -- [索引](00-index.md) — 文档导航与架构总览图 -- [系统总览](01-system-overview.md) — 分层架构与依赖注入 -- [协程池](02-worker-pool.md) — 有界并发与 per-user 限流 -- [RabbitMQ 集成](04-rabbitmq.md) — 持久化消息与重试机制 +- [[00-索引]] — 文档导航与架构总览图 +- [[01-系统总览]] — 分层架构与依赖注入 +- [[02-协程池]] — 有界并发与 per-user 限流 +- [[04-RabbitMQ集成]] — 持久化消息与重试机制 diff --git a/hzh/Gen2D/04-RabbitMQ集成.md b/hzh/Gen2D/04-RabbitMQ集成.md index cd621d6..bcba00e 100644 --- a/hzh/Gen2D/04-RabbitMQ集成.md +++ b/hzh/Gen2D/04-RabbitMQ集成.md @@ -1,17 +1,28 @@ -# 04 - RabbitMQ 集成 +--- +tags: [rabbitmq, amqp, message-queue, persistence, ack-nack, retry-pattern, go] +create time: 2026-06-03 10:15 +--- + +# 04. RabbitMQ 集成 + +## 概述 + +基于 AMQP 协议的持久化消息队列实现,支持手动 ACK/NACK、失败重试和死信丢弃,保障消息不丢失。 > **一句话概括**:基于 AMQP 的持久化消息队列,支持手动 ACK、失败重试和死信丢弃。 -## 消息流 +## 正文 + +### 消息流 ```mermaid graph TB subgraph Producer["生产者"] - HANDLER["Handler.Generate()"] + HANDLER["Handler.Generate"] end subgraph RabbitMQ["RabbitMQ"] - EXCHANGE["Default Exchange
(Direct)"] + EXCHANGE["Default ExchangeDirect"] QUEUE["gen2d:tasks
durable=true"] end @@ -19,15 +30,15 @@ graph TB CONSUME["channel.Consume
autoAck=false"] end - subgraph Decision["ACK/NACK 决策树"] - SUCCESS{"handler 成功?"} - RETRY{"retry < maxRetry?"} - ACK_OK["ACK
确认消费"] + subgraph Decision["ACK/NACK决策树"] + SUCCESS{"handler成功?"} + RETRY{"retry lt maxRetry?"} + ACK_OK["ACK确认消费"] NACK["NACK + requeue
重新入队"] - ACK_DISCARD["ACK (discard)
丢弃死信"] + ACK_DISCARD["ACK discard
丢弃死信"] end - HANDLER -->|"Publish
Persistent"| EXCHANGE + HANDLER -->|"Publish Persistent"| EXCHANGE EXCHANGE --> QUEUE QUEUE --> CONSUME CONSUME --> SUCCESS @@ -38,24 +49,16 @@ graph TB NACK -.->|"重新投递"| QUEUE ``` -## 连接流程 +### 连接流程 RabbitMQQueue 在初始化时完成连接、声明队列、设置 QoS: -``` -amqp.Dial(cfg.URL) - │ - ▼ -conn.Channel() - │ - ▼ -ch.QueueDeclare(name, durable=true, autoDelete=false, exclusive=false) - │ - ▼ -ch.Qos(prefetch=1, prefetchSize=0, global=false) - │ - ▼ -RabbitMQQueue{conn, channel, queue, maxRetry, prefetch} +```mermaid +graph TB + DIAL["amqp.Dialcfg.URL"] --> CHANNEL["conn.Channel"] + CHANNEL --> DECLARE["ch.QueueDeclare name durable=true"] + DECLARE --> QOS["ch.Qosprefetch=1"] + QOS --> RESULT["RabbitMQQueue实例"] ``` **参数说明**: @@ -67,9 +70,11 @@ RabbitMQQueue{conn, channel, queue, maxRetry, prefetch} | `Prefetch` | `1` | 每次预取消息数,1 保证公平调度 | | `MaxRetry` | `3` | 失败最大重试次数 | -> :bulb: **Prefetch=1 的含义**:每个 Consumer 同时只处理 1 条消息,处理完(ACK)后才接收下一条。这避免了消息堆积在 Consumer 端,配合协程池的并发控制实现精确的任务调度。 +> [!tip] Prefetch=1 的含义 +> +> 每个 Consumer 同时只处理 1 条消息,处理完ACK后才接收下一条。这避免了消息堆积在 Consumer 端,配合协程池的并发控制实现精确的任务调度。 -## 消息发布 (Submit) +### 消息发布 (Submit) ```go func (q *RabbitMQQueue) Submit(ctx context.Context, msg TaskMessage) error { @@ -100,7 +105,7 @@ func (q *RabbitMQQueue) Submit(ctx context.Context, msg TaskMessage) error { | `ContentType` | `application/json` | JSON 序列化 | | `x-retry-count` | `int` (header) | 当前重试次数,供消费端判断 | -## 消息消费 (Consume) +### 消息消费 (Consume) ```go func (q *RabbitMQQueue) Consume(ctx context.Context, handler func(TaskMessage) error) error { @@ -121,25 +126,25 @@ func (q *RabbitMQQueue) Consume(ctx context.Context, handler func(TaskMessage) e 手动 ACK 给予消费者完全的控制权——只有当消息被成功处理后才确认,否则可以选择重试或丢弃。 -## ACK/NACK 决策树 +### ACK/NACK 决策树 ```mermaid graph TB - MSG["收到消息"] --> PARSE{"JSON 解析成功?"} - PARSE -->|否| ACK_DISCARD1["ACK (discard)
格式错误无法恢复"] - PARSE -->|是| HANDLER{"handler(msg)
执行成功?"} - HANDLER -->|是| ACK_OK["ACK
确认消费"] - HANDLER -->|否| RETRY_CHECK{"msg.RetryCount
< maxRetry?"} - RETRY_CHECK -->|是| NACK["NACK(requeue=true)
重新入队等待重试"] - RETRY_CHECK -->|否| ACK_DISCARD2["ACK (discard)
超过最大重试,记录死信日志"] + MSG["收到消息"] --> PARSE{"JSON解析成功?"} + PARSE -->|否| ACK_DISCARD1["ACK discard
格式错误无法恢复"] + PARSE -->|是| HANDLER{"handlermsg执行成功?"} + HANDLER -->|是| ACK_OK["ACK确认消费"] + HANDLER -->|否| RETRY_CHECK{"msg.RetryCount lt maxRetry?"} + RETRY_CHECK -->|是| NACK["NACKrequeuetrue
重新入队等待重试"] + RETRY_CHECK -->|否| ACK_DISCARD2["ACK discard
超过最大重试记录死信日志"] ``` | 场景 | 操作 | 说明 | |------|------|------| -| handler 成功 | `d.Ack(false)` | 确认消费,消息从队列移除 | -| handler 失败 + retry < max | `d.Nack(false, true)` | 拒绝并重新入队,retry count 递增 | -| handler 失败 + retry >= max | `d.Ack(false)` + 日志 | 超过最大重试,丢弃(可扩展为死信队列) | -| JSON 解析失败 | `d.Ack(false)` | 格式错误无法恢复,直接丢弃 | +| handler 成功 | `d.Ackfalse` | 确认消费,消息从队列移除 | +| handler 失败 + retry < max | `d.Nackfalse, true` | 拒绝并重新入队,retry count 递增 | +| handler 失败 + retry >= max | `d.Ackfalse` + 日志 | 超过最大重试,丢弃可扩展为死信队列 | +| JSON 解析失败 | `d.Ackfalse` | 格式错误无法恢复,直接丢弃 | **重试计数传递**: @@ -155,9 +160,11 @@ if retry, ok := d.Headers["x-retry-count"].(int32); ok { } ``` -> :warning: **NACK requeue 的行为**:`Nack(false, true)` 会将消息重新放回队列头部。如果消费者立即再次消费,可能导致"毒消息"反复重试。Gen2D 通过 `maxRetry=3` 限制重试次数,并在超过后 ACK 丢弃来规避此问题。 +> [!warning] NACK requeue 的行为 +> +> `Nack(false, true)` 会将消息重新放回队列头部。如果消费者立即再次消费,可能导致"毒消息"反复重试。Gen2D 通过 `maxRetry=3` 限制重试次数,并在超过后 ACK 丢弃来规避此问题。 -## Metrics 指标 +### Metrics 指标 | Prometheus 指标 | 类型 | Label | 说明 | |-----------------|------|-------|------| @@ -167,7 +174,7 @@ if retry, ok := d.Headers["x-retry-count"].(int32); ok { | `gen2d_queue_errors_total` | Counter | `driver=rabbitmq`, `error_type` | 错误总量 | | `gen2d_queue_submit_duration_seconds` | Histogram | `driver=rabbitmq` | 发布耗时 | -## 关闭流程 +### 关闭流程 ```go func (q *RabbitMQQueue) Close() error { @@ -181,7 +188,7 @@ func (q *RabbitMQQueue) Close() error { **关闭顺序**:Channel 先于 Connection 关闭,确保所有未确认的消息被释放回队列。 -## 配置参考 +### 配置参考 ```yaml # config.yaml @@ -203,7 +210,7 @@ taskqueue: ## 关联文档 -- [索引](00-index.md) — 文档导航与架构总览图 -- [任务队列](03-task-queue.md) — 可插拔接口与 MemoryQueue 实现 -- [协程池](02-worker-pool.md) — 单机并发控制 -- [系统总览](01-system-overview.md) — 分层架构与配置级联 +- [[00-索引]] — 文档导航与架构总览图 +- [[03-任务队列]] — 可插拔接口与 MemoryQueue 实现 +- [[02-协程池]] — 单机并发控制 +- [[01-系统总览]] — 分层架构与配置级联 diff --git a/hzh/Gen2D/05-生成管线.md b/hzh/Gen2D/05-生成管线.md index b4e1de9..69b9fc8 100644 --- a/hzh/Gen2D/05-生成管线.md +++ b/hzh/Gen2D/05-生成管线.md @@ -1,6 +1,17 @@ -# 05 - 生成管线 (Generation Pipeline) +--- +tags: [pipeline, eino, graph-pattern, quality-check, fallback, state-machine, go] +create time: 2026-06-03 10:20 +--- -> **一句话概括**:基于 CloudWeGo Eino 的 4 阶段生成管线,带质量回退和降级,确保任务不阻塞。 +# 05. 生成管线 (Generation Pipeline) + +## 概述 + +基于 CloudWeGo Eino 的 4 阶段生成管线,带质量回退和降级,确保任务不阻塞。 + +--- + +## 正文 ## 管线拓扑 @@ -161,6 +172,15 @@ if imgCfg.APIKey == "" { 4. pass=false + RetryCount >= 3 → NextNode = "format_adapter"(降级) ``` +> [!tip] 重试策略思考 +> +> **为什么是 3 次?** +> - 第 1 次失败:LLM API 波动或偶发噪声,重试大概率通过 +> - 第 2 次失败:提示词可能不够精确,重新优化后改善 +> - 第 3 次仍失败:当前参数组合确实无法生成合格图片,继续重试只会浪费资源 +> +> 超过 3 次后 **降级到 FormatAdapter**——即使结果不完美,也比永远阻塞管线要好。这就是「宁可降级,不可阻塞」的原则。 + **路由分支**(Eino `AddBranch`): ```go @@ -185,7 +205,7 @@ g.AddBranch(nodeQualitySupervisor, compose.NewGraphBranch( | `NewCountedQualityChecker(n)` | 第 n 次调用后通过 | 测试重试逻辑 | | `AlwaysFailQualityChecker` | 始终返回 `false` | 测试降级路径 | -> :bulb: **可扩展性**:`QualityChecker` 是一个可替换的函数变量,未来可接入 LLM 视觉模型进行真正的质量评估。 +> [!tip] 可扩展性:`QualityChecker` 是一个可替换的函数变量,未来可接入 LLM 视觉模型进行真正的质量评估。 ### 4. FormatAdapter — 格式适配 @@ -318,7 +338,7 @@ type PipelineInput struct { ## 关联文档 -- [索引](00-index.md) — 文档导航与架构总览图 -- [系统总览](01-system-overview.md) — 分层架构与 Service 层定位 -- [协程池](02-worker-pool.md) — 管线执行的并发控制 -- [任务队列](03-task-queue.md) — 管线任务的排队机制 +- [[00-索引]] — 文档导航与架构总览图 +- [[01-系统总览]] — 分层架构与 Service 层定位 +- [[02-协程池]] — 管线执行的并发控制 +- [[03-任务队列]] — 管线任务的排队机制 diff --git a/hzh/Gen2D/06-精灵图处理.md b/hzh/Gen2D/06-精灵图处理.md index b79176f..e22bce8 100644 --- a/hzh/Gen2D/06-精灵图处理.md +++ b/hzh/Gen2D/06-精灵图处理.md @@ -1,18 +1,27 @@ -# 06 — 精灵图处理管线 +--- +tags: [image-processing, sprite-sheet, gif, computer-vision, go, algorithm] +create time: 2026-06-03 10:25 +--- -> **一句话概括**:自动精灵图处理管线 — 背景移除 → 投影检测 → 切割 → 对齐 → GIF 预览,将一张 AI 生成的 Sprite Sheet 无缝转化为可用的逐帧动画资源。 +# 06. 精灵图处理管线 + +## 概述 + +自动精灵图处理管线 — 背景移除 → 投影检测 → 切割 → 对齐 → GIF 预览,将一张 AI 生成的 Sprite Sheet 无缝转化为可用的逐帧动画资源。 --- +## 正文 + ```mermaid flowchart LR - A["🖼️ Input PNG
Sprite Sheet"] --> B["🪄 Background
Removal"] - B --> C["📊 Gap
Detection"] - C --> D["✂️ Tile
Extract"] - D --> E["🔍 Filter
MinFill"] - E --> F["📐 Trim
Alpha"] - F --> G["🎯 Align
padToLargest"] - G --> H["🎬 GIF
Preview"] + A["Input PNG"] --> B["Background Removal"] + B --> C["Gap Detection"] + C --> D["Tile Extract"] + D --> E["Filter MinFill"] + E --> F["Trim Alpha"] + F --> G["Align padToLargest"] + G --> H["GIF Preview"] style A fill:#e3f2fd,stroke:#1976d2 style B fill:#fff3e0,stroke:#f57c00 @@ -26,7 +35,7 @@ flowchart LR --- -## 📋 处理管线总览 +## 处理管线总览 AI 生成的 Sprite Sheet 通常包含白色/绿色背景、不均匀间距和尺寸差异。Gen2D 的精灵图处理管线自动完成从"原始 PNG"到"可用动画帧"的全部转换工作,无需用户手动操作。 @@ -42,7 +51,7 @@ AI 生成的 Sprite Sheet 通常包含白色/绿色背景、不均匀间距和 --- -## 🪄 步骤 1:背景移除 +## 步骤 1:背景移除 AI 生成的图片通常带有纯色背景,管线支持两种模式: @@ -69,11 +78,12 @@ if gDominance > tolerance*255 → 透明化 - **GreenTolerance** 默认 0.2,控制绿色检测灵敏度 - 适用于绿色背景的 AI 生成图 -> 💡 **设计选择**:两种模式互斥,`Process()` 根据 `opts.WhiteBg` / `opts.GreenScreen` 自动选择。 +> [!tip] 设计选择 +> 两种模式互斥,`Process()` 根据 `opts.WhiteBg` / `opts.GreenScreen` 自动选择。 --- -## 📊 步骤 2:切割策略 +## 步骤 2:切割策略 管线提供两种切割方式,根据配置自动切换: @@ -90,6 +100,10 @@ if gDominance > tolerance*255 → 透明化 - 防止突出物(武器/尾巴)被其他行稀释 - 每行段获得独立的列边界,互不干扰 +> [!question] 为什么不用简单的阈值分割? +> +> 精灵图中的角色往往有复杂轮廓——比如挥舞的剑可能横跨多个帧的位置。如果仅用固定阈值,剑的连续像素会让算法误判为一帧。**按行段独立分析**的思路是把二维问题拆解为多个一维子问题,每个子问题只关心当前行段的内容,从而避免跨行干扰。 + ```mermaid flowchart TB A["输入图像"] --> B["全局行投影
检测行间隙"] @@ -124,7 +138,7 @@ flowchart TB --- -## 🔍 步骤 3:过滤 — MinFillRatio +## 步骤 3:过滤 — MinFillRatio 切割后的每个 tile 都计算填充率: @@ -138,7 +152,7 @@ fillRatio = 非透明像素数 / 总像素数 --- -## 📐 步骤 4:裁剪 — trimAlpha +## 步骤 4:裁剪 — trimAlpha 对每个 tile 执行透明边框裁剪: @@ -148,7 +162,7 @@ fillRatio = 非透明像素数 / 总像素数 --- -## 🎯 步骤 5:对齐 — padToLargest +## 步骤 5:对齐 — padToLargest 动画播放时,如果每帧尺寸不同且内容未对齐,会导致角色"抖动"。 @@ -169,7 +183,7 @@ canvasH = maxH * 110% // 最大帧高度 + 10% padding --- -## 🎬 GIF Maker +## GIF Maker `gifmaker.Encode()` 将处理后的帧序列编码为动画 GIF: @@ -186,11 +200,12 @@ anim.Disposal = append(anim.Disposal, gif.DisposalBackground) anim.BackgroundIndex = 0 // 透明色 ``` -> ⚠️ **DisposalBackground 的重要性**:如果不设置此选项,GIF 播放器会在前一帧基础上叠加新帧,产生"残影"效果。 +> [!warning] DisposalBackground 的重要性 +> 如果不设置此选项,GIF 播放器会在前一帧基础上叠加新帧,产生"残影"效果。 --- -## 📦 Options 配置速查 +## Options 配置速查 | 参数 | 类型 | 默认值 | 说明 | |------|------|--------|------| @@ -207,8 +222,8 @@ anim.BackgroundIndex = 0 // 透明色 --- -## 🔗 关联文档 +## 关联文档 -- [← 返回索引](00-index.md) -- [05 — 生成管线](05-generation-pipeline.md) — 管线中 SplitSprite 节点的调用方 -- [08 — SSE 实时推送](08-sse-push.md) — 处理进度的实时推送 +- [[00-索引]] — 文档导航与架构总览图 +- [[05-生成管线]] — 管线中 SplitSprite 节点的调用方 +- [[08-SSE实时推送]] — 处理进度的实时推送 diff --git a/hzh/Gen2D/07-可观测性.md b/hzh/Gen2D/07-可观测性.md index b0bb690..e68d2a2 100644 --- a/hzh/Gen2D/07-可观测性.md +++ b/hzh/Gen2D/07-可观测性.md @@ -1,28 +1,19 @@ -# 07 — 可观测性 +--- +tags: [observability, prometheus, grafana, metrics, alerting, monitoring] +create time: 2026-06-03 10:30 +--- -> **一句话概括**:35 个 Prometheus 指标 + 3 个 Grafana 仪表盘 + 10 条告警规则,覆盖全栈,让系统运行状态一目了然。 +# 07. 可观测性 + +## 概述 + +35 个 Prometheus 指标 + 3 个 Grafana 仪表盘 + 10 条告警规则,覆盖全栈,让系统运行状态一目了然。 --- -```mermaid -flowchart LR - A["🌐 Request"] --> B["📊 Metrics
Middleware"] - B --> C["⚙️ Application
Logic"] - C --> D["📈 Prometheus
Scrape"] - D --> E["📉 Grafana
Dashboard"] - D --> F["🚨 AlertManager
Notify"] +## 正文 - style A fill:#e3f2fd,stroke:#1976d2 - style B fill:#fff3e0,stroke:#f57c00 - style C fill:#e8f5e9,stroke:#388e3c - style D fill:#fce4ec,stroke:#c62828 - style E fill:#e8eaf6,stroke:#303f9f - style F fill:#ffebee,stroke:#b71c1c -``` - ---- - -## 📐 指标体系总览 +## 指标体系总览 Gen2D 遵循 Prometheus 命名最佳实践,所有指标使用 `gen2d_` 前缀,共 **35 个指标**,分为 **6 大组**: @@ -37,9 +28,9 @@ Gen2D 遵循 Prometheus 命名最佳实践,所有指标使用 `gen2d_` 前缀 --- -## 📊 五大指标组详解 +## 五大指标组详解 -### 1️⃣ HTTP 层指标 +### HTTP 层指标 Gin 中间件自动采集,**零业务代码侵入**。 @@ -51,9 +42,10 @@ Gin 中间件自动采集,**零业务代码侵入**。 | `gen2d_http_response_size_bytes` | Histogram | method, path | 响应体大小 | | `gen2d_http_requests_in_flight` | Gauge | — | 当前并发请求数 | -> 💡 **FullPath() 的关键作用**:使用路由模板 `/api/v1/tasks/:taskId` 而非实际路径 `/api/v1/tasks/abc123`,避免高基数标签导致 Prometheus 内存爆炸。 +> [!tip] FullPath() 的关键作用 +> 使用路由模板 `/api/v1/tasks/:taskId` 而非实际路径 `/api/v1/tasks/abc123`,避免高基数标签导致 Prometheus 内存爆炸。 -### 2️⃣ 限流层指标 +### 限流层指标 | 指标名 | 类型 | 标签 | 说明 | |--------|------|------|------| @@ -63,7 +55,7 @@ Gin 中间件自动采集,**零业务代码侵入**。 - `scope`:`user` / `global` - `result`:`allowed` / `denied` -### 3️⃣ 任务队列层指标 +### 任务队列层指标 | 指标名 | 类型 | 标签 | 说明 | |--------|------|------|------| @@ -75,7 +67,7 @@ Gin 中间件自动采集,**零业务代码侵入**。 - `driver`:`memory` / `rabbitmq` -### 4️⃣ 协程池层指标 +### 协程池层指标 | 指标名 | 类型 | 标签 | 说明 | |--------|------|------|------| @@ -86,7 +78,7 @@ Gin 中间件自动采集,**零业务代码侵入**。 | `gen2d_pool_rejected_total` | Counter | — | 被拒绝的任务 | | `gen2d_pool_task_duration_seconds` | Histogram | — | 任务执行耗时 | -### 5️⃣ Pipeline 业务层指标 +### Pipeline 业务层指标 | 指标名 | 类型 | 标签 | 说明 | |--------|------|------|------| @@ -96,9 +88,10 @@ Gin 中间件自动采集,**零业务代码侵入**。 | `gen2d_pipeline_retries_total` | CounterVec | stage | 各阶段重试次数 | | `gen2d_pipeline_tasks_active` | Gauge | — | 当前执行中的 Pipeline 数 | -> 🔍 **stage_duration 定位瓶颈**:通过 `stage` 标签(如 `asset_generator`、`quality_check`)可以精确定位哪个阶段是性能瓶颈。 +> [!tip] stage_duration 定位瓶颈 +> 通过 `stage` 标签(如 `asset_generator`、`quality_check`)可以精确定位哪个阶段是性能瓶颈。 -### 6️⃣ 基础设施层指标 +### 基础设施层指标 | 指标名 | 类型 | 标签 | 说明 | |--------|------|------|------| @@ -110,7 +103,7 @@ Gin 中间件自动采集,**零业务代码侵入**。 --- -## 🔧 中间件集成 +## 中间件集成 ### Metrics 中间件工作流程 @@ -136,7 +129,7 @@ func Metrics() gin.HandlerFunc { --- -## 📉 Grafana 仪表盘 +## Grafana 仪表盘 | 仪表盘 | 用途 | 关键面板 | |--------|------|---------| @@ -146,7 +139,7 @@ func Metrics() gin.HandlerFunc { --- -## 🚨 告警规则 +## 告警规则 共 **10 条告警规则**,覆盖限流、队列、协程池、管线和基础设施: @@ -163,11 +156,12 @@ func Metrics() gin.HandlerFunc { | `HighErrorRate` | 🔴 | 5xx 错误率 > 5%(持续 5m) | 服务异常 | | `HighLatency` | ⚠️ | P95 延迟 > 5s(持续 5m) | 影响用户体验 | -> 🛡️ **告警级别说明**:🔴 Critical 表示需要立即处理,⚠️ Warning 表示需要关注但不紧急。 +> [!note] 告警级别说明 +> Critical 表示需要立即处理,Warning 表示需要关注但不紧急。 --- -## 📦 基础设施指标 +## 基础设施指标 除业务指标外,Gen2D 还监控外部依赖的健康状态: @@ -194,9 +188,8 @@ flowchart LR --- -## 🔗 关联文档 +## 关联文档 -- [← 返回索引](00-index.md) -- [10 — 中间件链](10-middleware-chain.md) — Metrics 中间件的挂载位置 -- [09 — 限流](09-rate-limiting.md) — 限流指标的采集方式 -- [14 — 部署架构](14-deployment.md) — Prometheus + Grafana 的部署配置 +- [[10-中间件链]] — Metrics 中间件的挂载位置 +- [[09-限流]] — 限流指标的采集方式 +- [[14-部署架构]] — Prometheus + Grafana 的部署配置 diff --git a/hzh/Gen2D/08-SSE实时推送.md b/hzh/Gen2D/08-SSE实时推送.md index 973d53f..dc5e189 100644 --- a/hzh/Gen2D/08-SSE实时推送.md +++ b/hzh/Gen2D/08-SSE实时推送.md @@ -1,16 +1,25 @@ -# 08 — SSE 实时推送 +--- +tags: [sse, event-stream, pub-sub, real-time, websockets-alternative, go] +create time: 2026-06-03 10:35 +--- -> **一句话概括**:内存 EventBus 发布/订阅,SSE 推送管线进度到浏览器,让用户实时看到生成过程。 +# 08. SSE 实时推送 + +## 概述 + +内存 EventBus 发布/订阅,SSE 推送管线进度到浏览器,让用户实时看到生成过程。 --- +## 正文 + ```mermaid flowchart LR - P["⚙️ Pipeline
Callback"] -->|"Publish"| EB["📡 EventBus
Broker"] - EB -->|"Subscribe
taskID"| H1["🌐 SSE Handler
/tasks/:id/stream"] - EB -->|"SubscribeAll
global"| H2["🌐 SSE Handler
/projects/:id/stream"] - H1 -->|"text/event-stream"| B1["🖥️ Browser
EventSource"] - H2 -->|"text/event-stream"| B2["🖥️ Browser
EventSource"] + P["Pipeline Callback"] -->|"Publish"| EB["EventBus Broker"] + EB -->|"Subscribe taskID"| H1["SSE Handler /tasks/:id/stream"] + EB -->|"SubscribeAll global"| H2["SSE Handler /projects/:id/stream"] + H1 -->|"text/event-stream"| B1["Browser EventSource"] + H2 -->|"text/event-stream"| B2["Browser EventSource"] style P fill:#e8f5e9,stroke:#388e3c style EB fill:#fff3e0,stroke:#f57c00 @@ -22,7 +31,7 @@ flowchart LR --- -## 📡 EventBus 架构 +## EventBus 架构 EventBus 是 Gen2D 的内存事件总线,负责在 Pipeline 执行过程中发布进度事件,并由 SSE Handler 订阅推送给客户端。 @@ -57,7 +66,7 @@ type TaskEvent struct { --- -## 🔔 订阅模式 +## 订阅模式 ### Subscribe(taskID) — 任务级订阅 @@ -93,7 +102,7 @@ func (b *Broker) SubscribeAll() <-chan TaskEvent { --- -## 📤 Publish — 扇出分发 +## Publish — 扇出分发 ```go func (b *Broker) Publish(taskID string, event TaskEvent) { @@ -119,7 +128,7 @@ func (b *Broker) Publish(taskID string, event TaskEvent) { --- -## 🌐 SSE Handler +## SSE Handler ### Stream — 任务级流 @@ -175,7 +184,7 @@ GET /api/v1/projects/:projectId/stream --- -## 📊 数据流全景 +## 数据流全景 ```mermaid sequenceDiagram @@ -199,7 +208,7 @@ sequenceDiagram --- -## 🛡️ 容错设计 +## 容错设计 | 场景 | 处理方式 | |------|---------| @@ -211,9 +220,8 @@ sequenceDiagram --- -## 🔗 关联文档 +## 关联文档 -- [← 返回索引](00-index.md) -- [05 — 生成管线](05-generation-pipeline.md) — Pipeline 中的进度回调 -- [07 — 可观测性](07-observability.md) — SSE 连接的监控 -- [10 — 中间件链](10-middleware-chain.md) — SSE 端点的中间件配置 +- [[05-生成管线]] — Pipeline 中的进度回调 +- [[07-可观测性]] — SSE 连接的监控 +- [[10-中间件链]] — SSE 端点的中间件配置 diff --git a/hzh/Gen2D/09-限流.md b/hzh/Gen2D/09-限流.md index 68d80fc..492dcd9 100644 --- a/hzh/Gen2D/09-限流.md +++ b/hzh/Gen2D/09-限流.md @@ -1,16 +1,25 @@ -# 09 — 限流 +--- +tags: [rate-limiting, redis, lua, token-bucket, distributed-system, go] +create time: 2026-06-03 10:40 +--- -> **一句话概括**:Redis Lua 原子令牌桶 + 双层限流 + Fail-Open 降级,保护系统免受过载。 +# 09. 限流 + +## 概述 + +Redis Lua 原子令牌桶 + 双层限流 + Fail-Open 降级,保护系统免受过载。 --- +## 正文 + ```mermaid flowchart LR - A["🌐 Request"] --> B["🌍 Global
Limiter"] - B -->|pass| C["👤 User
Limiter"] - B -->|deny| F["❌ 429"] + A["Request"] --> B["Global Limiter"] + B -->|pass| C["User Limiter"] + B -->|deny| F["429"] B -->|redis-fail| C - C -->|pass| D["✅ Handler"] + C -->|pass| D["Handler"] C -->|deny| F C -->|redis-fail| D @@ -23,7 +32,7 @@ flowchart LR --- -## ⚙️ 令牌桶算法 +## 令牌桶算法 Gen2D 使用 **Redis + Lua 脚本** 实现分布式令牌桶限流,保证原子性和一致性。 @@ -75,11 +84,12 @@ return {allowed, tokens, retry_after} - 用完后不补充(`rate = 0` 时跳过 refill) - 等待 key 过期后重置(`Expiration` 控制窗口大小) -> 💡 **适用场景**:24 小时维度的配额控制,如"每天 30 次提示词优化"。 +> [!tip] 适用场景 +> 24 小时维度的配额控制,如"每天 30 次提示词优化"。 --- -## 🔀 双层限流配置 +## 双层限流配置 Gen2D 对核心接口实施**全局限流 + 用户限流**双重保护: @@ -126,7 +136,7 @@ type Config struct { --- -## 🛡️ Fail-Open 降级 +## Fail-Open 降级 当 Redis 不可用时限流器自动降级为 **Fail-Open** 模式: @@ -148,11 +158,20 @@ func (l *TokenBucketLimiter) Allow(ctx context.Context, key string) (bool, int, | **Fail-Open** ✅ | 保证可用性,用户体验不受影响 | 可能短暂失去限流保护 | | Fail-Close | 严格限流保护 | Redis 故障导致全站不可用 | -> 🛡️ **选择 Fail-Open**:在"偶尔超限"和"完全不可用"之间,优先保证服务可用性。 +> [!question] 为什么选择 Fail-Open 而不是 Fail-Close? +> +> 这是 **「可用性 vs 安全性」** 的经典抉择。在 Gen2D 的场景中: +> - 限流失效的代价:短时间内有人可能超出配额(几分钟到几小时) +> - 限流强固化的代价:**所有用户都无法使用服务** +> +> 显然,前者是可以接受的风险——超出配额的用户可以后续通过账单追缴;而后者意味着业务完全停摆。这种「宁可放宽、不可收紧」的设计哲学在基础设施层非常重要。 + +> [!note] 选择 Fail-Open +> 在"偶尔超限"和"完全不可用"之间,优先保证服务可用性。 --- -## 📡 中间件响应 +## 中间件响应 限流中间件返回标准化的 HTTP 响应: @@ -183,7 +202,7 @@ X-RateLimit-Remaining: 0 --- -## 📊 指标采集 +## 指标采集 限流中间件自动采集 Prometheus 指标: @@ -200,8 +219,7 @@ metrics.RateLimitRemainingTokens.WithLabelValues(scope, endpoint).Set(float64(re --- -## 🔗 关联文档 +## 关联文档 -- [← 返回索引](00-index.md) -- [10 — 中间件链](10-middleware-chain.md) — 限流中间件在链中的位置 -- [07 — 可观测性](07-observability.md) — 限流指标和告警规则 +- [[10-中间件链]] — 限流中间件在链中的位置 +- [[07-可观测性]] — 限流指标和告警规则 diff --git a/hzh/Gen2D/10-中间件链.md b/hzh/Gen2D/10-中间件链.md index 28baae3..4ae93d1 100644 --- a/hzh/Gen2D/10-中间件链.md +++ b/hzh/Gen2D/10-中间件链.md @@ -1,23 +1,32 @@ -# 10 — 中间件链 - -> **一句话概括**:Logger → Recovery → Metrics → Auth → RateLimit → Handler,洋葱模型,层层守护请求处理。 - --- +tags: [middleware, gin, onion-pattern, auth, logging, recovery] +create time: 2026-06-03 10:45 +--- + +# 10. 中间件链 + +## 概述 + +Logger → Recovery → Metrics → Auth → RateLimit → Handler,洋葱模型,层层守护请求处理。 + +## 正文 + +### 请求处理链路 ```mermaid flowchart LR - A["🌐 Request"] --> B["📝 Logger"] - B --> C["🛡️ Recovery"] - C --> D["📊 Metrics"] - D --> E["🔐 Auth"] - E --> F["🚦 RateLimit"] - F --> G["⚙️ Handler"] + A["Request"] --> B["Logger"] + B --> C["Recovery"] + C --> D["Metrics"] + D --> E["Auth"] + E --> F["RateLimit"] + F --> G["Handler"] G --> F F --> E E --> D D --> C C --> B - B --> H["📡 Response"] + B --> H["Response"] style A fill:#e3f2fd,stroke:#1976d2 style B fill:#e8f5e9,stroke:#388e3c @@ -29,24 +38,17 @@ flowchart LR style H fill:#e8eaf6,stroke:#303f9f ``` ---- - -## 🧅 洋葱模型 +### 洋葱模型 Gin 的中间件采用**洋葱模型**:请求从外到内穿过各中间件,响应从内到外返回。每个中间件可以在 `c.Next()` 前后执行逻辑。 -``` -请求 → Logger.enter → Recovery.enter → Metrics.enter → Auth.enter → RateLimit.enter → Handler -响应 ← Logger.leave ← Recovery.leave ← Metrics.leave ← Auth.leave ← RateLimit.leave ← Handler -``` - --- -## 📝 Logger — 请求日志 +### Logger — 请求日志 **职责**:为每个请求生成唯一 ID,记录请求详情。 -### 核心逻辑 +#### 核心逻辑 ```go func Logger() gin.HandlerFunc { @@ -63,7 +65,7 @@ func Logger() gin.HandlerFunc { } ``` -### RequestID 生成 +#### RequestID 生成 ```go func generateRequestID() string { @@ -80,7 +82,7 @@ func generateRequestID() string { | 时间戳(毫秒) | 保证时间有序性 | | 8 字节随机 hex | 保证唯一性 | -### 日志级别映射 +#### 日志级别映射 | HTTP 状态码 | 日志级别 | 含义 | |:-----------:|:-------:|------| @@ -90,7 +92,7 @@ func generateRequestID() string { --- -## 🛡️ Recovery — Panic 恢复 +### Recovery — Panic 恢复 **职责**:捕获未处理的 panic,防止服务崩溃。 @@ -124,7 +126,7 @@ func Recovery() gin.HandlerFunc { --- -## 📊 Metrics — 指标采集 +### Metrics — 指标采集 **职责**:自动采集 HTTP 请求的性能指标。 @@ -152,11 +154,12 @@ func Metrics() gin.HandlerFunc { | ResponseSize | `c.Next()` 之后 | Writer 此时已写入 | | Duration | `c.Next()` 之后 | 需要计算总耗时 | -> 💡 **FullPath() 的关键作用**:使用路由模板 `/api/v1/tasks/:taskId` 而非实际路径,避免高基数标签导致 Prometheus 内存爆炸。 +> [!tip] FullPath() 的关键作用 +> 使用路由模板 `/api/v1/tasks/:taskId` 而非实际路径,避免高基数标签导致 Prometheus 内存爆炸。 --- -## 🔐 Auth — JWT 认证 +### Auth — JWT 认证 **职责**:验证 Bearer token,提取用户身份。 @@ -188,13 +191,13 @@ func AuthMiddleware(jwtSecret string) gin.HandlerFunc { | 场景 | HTTP 状态码 | 消息 | |------|:-----------:|------| | 无 token | 401 | 未提供认证令牌 | -| 格式错误 | 401 | 认证格式错误,需为 Bearer \ | +| 格式错误 | 401 | 认证格式错误,需为 Bearer | | token 无效 | 401 | 令牌无效或已过期 | | 解析失败 | 401 | 令牌解析失败 | --- -## 🚦 RateLimit — 限流 +### RateLimit — 限流 **职责**:按路由配置执行双层限流(全局 + 用户)。 @@ -207,13 +210,13 @@ v1Auth.POST("/generate", ) ``` -详见 [09 — 限流](09-rate-limiting.md)。 +详见 [[09-限流]]。 --- -## 🔗 中间件挂载 +### 中间件挂载 -### 全局链(所有请求) +#### 全局链(所有请求) ```go r := gin.New() @@ -222,14 +225,14 @@ r.Use(mildware.Recovery()) // 2. Panic 恢复 r.Use(mildware.Metrics()) // 3. 指标采集 ``` -### 路由组级链(需认证) +#### 路由组级链(需认证) ```go v1Auth := r.Group("/api/v1") v1Auth.Use(mildware.AuthMiddleware(cfg.JWT.Secret)) // 4. JWT 认证 ``` -### 端点级链(需限流) +#### 端点级链(需限流) ```go v1Auth.POST("/generate", @@ -241,7 +244,7 @@ v1Auth.POST("/generate", --- -## 📋 端点链示例:/api/v1/generate +### 端点链示例:/api/v1/generate ```mermaid flowchart TB @@ -255,7 +258,7 @@ flowchart TB H --> I["Metrics
记录 duration, respSize"] I --> J["Recovery
检查是否 panic"] J --> K["Logger
记录请求日志"] - K --> L["📡 Response"] + K --> L["Response"] style A fill:#e3f2fd,stroke:#1976d2 style B fill:#e8f5e9,stroke:#388e3c @@ -268,7 +271,7 @@ flowchart TB style L fill:#e8eaf6,stroke:#303f9f ``` -### 完整请求生命周期 +#### 完整请求生命周期 | 阶段 | 中间件 | 动作 | |:----:|--------|------| @@ -285,7 +288,7 @@ flowchart TB --- -## 📊 中间件职责矩阵 +### 中间件职责矩阵 | 中间件 | 请求进入 | 请求离开 | 异常处理 | 作用范围 | |--------|---------|---------|---------|---------| @@ -297,9 +300,8 @@ flowchart TB --- -## 🔗 关联文档 +## 关联文档 -- [← 返回索引](00-index.md) -- [07 — 可观测性](07-observability.md) — Metrics 中间件采集的指标 -- [09 — 限流](09-rate-limiting.md) — RateLimit 中间件的详细实现 -- [08 — SSE 实时推送](08-sse-push.md) — SSE 端点的中间件配置 +- [[07-可观测性]] — Metrics 中间件采集的指标 +- [[09-限流]] — RateLimit 中间件的详细实现 +- [[08-SSE实时推送]] — SSE 端点的中间件配置 diff --git a/hzh/Gen2D/11-Consumer-Producer桥接.md b/hzh/Gen2D/11-Consumer-Producer桥接.md index 7f82b64..2174c26 100644 --- a/hzh/Gen2D/11-Consumer-Producer桥接.md +++ b/hzh/Gen2D/11-Consumer-Producer桥接.md @@ -1,34 +1,41 @@ -# 11. Consumer-Producer 桥接模式 - -> **一句话概括**:`Consumer` 结构体桥接 `TaskQueue` 和 `WorkerPool`,实现生产者与消费者的彻底解耦。 - +--- +tags: [consumer-producer, bridge-pattern, decoupling, go, message-queue] +create time: 2026-06-03 10:50 --- -## 架构总览 +# 11. Consumer-Producer 桥接模式 + +## 概述 + +Consumer 结构体桥接 TaskQueue 和 WorkerPool,实现生产者与消费者的彻底解耦。 + +## 正文 + +### 架构总览 ```mermaid flowchart LR - subgraph Producer["生产者"] - A[Generate Handler] + subgraph Producer["Producer"] + A["Generate Handler"] end subgraph Queue["TaskQueue 接口"] - B((Memory\nQueue)) - C((RabbitMQ\nQueue)) + B[("Memory Queue")] + C[("RabbitMQ Queue")] end subgraph Bridge["Consumer 桥接层"] - D{{"Consumer\n(bridge)"}} + D{{"Consumer"}} end subgraph Pool["WorkerPool"] - E[Worker 1] - F[Worker 2] - G[Worker N] + E["Worker 1"] + F["Worker 2"] + G["Worker N"] end subgraph Pipeline["业务逻辑"] - H[Eino Pipeline] + H["Eino Pipeline"] end A -->|"Submit(msg)"| B @@ -45,9 +52,7 @@ flowchart LR style D fill:#f9a825,stroke:#333,color:#000 ``` ---- - -## 核心结构体 +### 核心结构体 `Consumer` 是整个任务调度体系的**桥梁**,它只做一件事:从队列取消息,提交到协程池。 @@ -69,14 +74,14 @@ type Consumer struct { --- -## 工作流程 +### 工作流程 -### 启动消费循环 +#### 启动消费循环 ```mermaid sequenceDiagram participant Main as main.go - participant Consumer + participant Consumer as Consumer participant Queue as TaskQueue participant Pool as WorkerPool participant Handler as RunFromTaskMessage @@ -102,20 +107,20 @@ sequenceDiagram --- -## 解耦的三层设计 +### 解耦的三层设计 ```mermaid flowchart TB - subgraph "第 1 层:消息源" + subgraph "消息源" Q["TaskQueue 接口\n(Memory / RabbitMQ)"] end - subgraph "第 2 层:桥接" + subgraph "桥接" C["Consumer\n(只关心 消费→提交)"] end - subgraph "第 3 层:执行引擎" + subgraph "执行引擎" P["WorkerPool\n(只关心 并发控制)"] end - subgraph "第 4 层:业务逻辑" + subgraph "业务逻辑" H["TaskHandler 回调\n(RunFromTaskMessage)"] end @@ -133,7 +138,7 @@ flowchart TB --- -## Handler 层注入 +### Handler 层注入 `Consumer` 不硬编码业务逻辑,而是通过 `TaskHandler` 函数签名由外部注入: @@ -149,13 +154,13 @@ consumer := worker.NewConsumer(tq, pool, handler.RunFromTaskMessage) --- -## 信号处理与优雅关闭 +### 信号处理与优雅关闭 ```mermaid sequenceDiagram - participant OS as 操作系统 + participant OS as OS participant Main as main.go - participant Consumer + participant Consumer as Consumer participant Pool as WorkerPool OS->>Main: SIGINT / SIGTERM @@ -173,13 +178,9 @@ sequenceDiagram 2. **再关 WorkerPool** — 等待已提交的任务执行完毕(最多 30 秒) 3. **最后关闭队列连接** — 释放 RabbitMQ / 内存资源 ---- - ## 关联文档 -| 文档 | 关系 | -|------|------| -| [03 - 任务队列](03-task-queue.md) | Consumer 的消息来源 | -| [02 - 协程池](02-worker-pool.md) | Consumer 的执行引擎 | -| [12 - 三级降级策略](12-three-tier-fallback.md) | Consumer 不参与降级,降级在 Handler 层 | -| [00 - 索引](00-index.md) | 返回文档总览 | +- [[03-任务队列]] — Consumer 的消息来源 +- [[02-协程池]] — Consumer 的执行引擎 +- [[12-三级降级策略]] — Consumer 不参与降级,降级在 Handler 层 +- [[00-索引]] — 返回文档总览 diff --git a/hzh/Gen2D/12-三级降级策略.md b/hzh/Gen2D/12-三级降级策略.md index ff21bc0..86b8946 100644 --- a/hzh/Gen2D/12-三级降级策略.md +++ b/hzh/Gen2D/12-三级降级策略.md @@ -1,25 +1,32 @@ -# 12. 三级降级策略 - -> **一句话概括**:Generate handler 实现三级降级 — TaskQueue -> WorkerPool -> Legacy FIFO,宁可降级也不拒绝服务。 - +--- +tags: [fallback-pattern, resilience, degradation, error-handling, go] +create time: 2026-06-03 10:55 --- -## 降级链总览 +# 12. 三级降级策略 + +## 概述 + +Generate handler 实现三级降级 — TaskQueue -> WorkerPool -> Legacy FIFO,宁可降级也不拒绝服务。 + +## 正文 + +### 降级链总览 ```mermaid flowchart TD - R["HTTP Request\nPOST /api/v1/generate"] --> A{"TaskQueue\n可用?"} + R["HTTP Request\nPOST /api/v1/generate"] --> A{"TaskQueue 可用?"} A -->|"Submit 成功"| S1["200 OK\ntaskId 返回"] - A -->|"Submit 失败"| B{"WorkerPool\n可用?"} + A -->|"Submit 失败"| B{"WorkerPool 可用?"} B -->|"Submit 成功"| S2["200 OK\ntaskId 返回"] B -->|"ErrPoolFull"| E2["503 Service\nUnavailable"] - B -->|"ErrUserLimit"| E3["429 Too Many\nRequests"] - B -->|"Pool 不可用"| C{"Legacy FIFO\nQueue 可用?"} + B -->|"ErrUserLimit"| E3["429 Too Many Requests"] + B -->|"Pool 不可用"| C{"Legacy FIFO Queue 可用?"} C -->|"Enqueue 成功"| S3["200 OK\ntaskId 返回"] - C -->|"Queue 不可用"| E4["500 Internal\nServer Error"] + C -->|"Queue 不可用"| E4["500 Internal Server Error"] style A fill:#4caf50,stroke:#333,color:#fff style B fill:#ff9800,stroke:#333,color:#fff @@ -29,11 +36,9 @@ flowchart TD style S3 fill:#8bc34a,stroke:#333,color:#fff ``` ---- +### 三级详解 -## 三级详解 - -### 第一级:TaskQueue(优先路径) +#### 第一级:TaskQueue(优先路径) | 属性 | 说明 | |------|------| @@ -55,7 +60,7 @@ if taskQueue != nil { } ``` -### 第二级:WorkerPool(有界并发) +#### 第二级:WorkerPool(有界并发) | 属性 | 说明 | |------|------| @@ -80,7 +85,7 @@ if workerPool != nil { } ``` -### 第三级:Legacy FIFO(最后保底) +#### 第三级:Legacy FIFO(最后保底) | 属性 | 说明 | |------|------| @@ -98,7 +103,7 @@ c.JSON(200, taskId) --- -## HTTP 状态码映射 +### HTTP 状态码映射 ```mermaid flowchart LR @@ -133,7 +138,7 @@ flowchart LR --- -## 为什么需要三级? +### 为什么需要三级? ```mermaid flowchart TB @@ -157,9 +162,18 @@ flowchart TB | WorkerPool | 并发控制 + 背压 | 单机部署,需要限制资源 | | Legacy FIFO | 可用性兜底 | 开发/测试环境,或队列组件故障 | +> [!tip] 分级降级的核心原则 +> +> **每一级都是前一级功能的超集**。也就是说: +> - TaskQueue = 持久化排队 + 重试 + Consumer 消费 + WorkerPool 执行 +> - WorkerPool = 有界并发 + 直接执行(跳过持久化和消费者) +> - Legacy FIFO = 最简串行队列(跳过所有高级特性) +> +> 这种设计确保降级过程是**渐进的**——功能逐步减少但服务始终可用。类比现实中的「应急灯」:市电断了 → 应急灯亮 → 最差情况还有手电筒。永远保留一条最低限度的通路。 + --- -## 设计哲学 +### 设计哲学 > **宁可降级,也不能拒绝服务。** @@ -171,13 +185,9 @@ flowchart TB 每一级降级都意味着功能的减少,但服务的可用性始终得到保障。前端收到 `200 OK` 后,通过 SSE 或轮询获取任务进度,对用户而言体验一致。 ---- - ## 关联文档 -| 文档 | 关系 | -|------|------| -| [11 - Consumer-Producer 桥接](11-consumer-producer.md) | Consumer 连接 TaskQueue 和 WorkerPool | -| [02 - 协程池](02-worker-pool.md) | WorkerPool 的背压和限流机制 | -| [03 - 任务队列](03-task-queue.md) | TaskQueue 接口的可插拔设计 | -| [00 - 索引](00-index.md) | 返回文档总览 | +- [[11-Consumer-Producer桥接]] — Consumer 连接 TaskQueue 和 WorkerPool +- [[02-协程池]] — WorkerPool 的背压和限流机制 +- [[03-任务队列]] — TaskQueue 接口的可插拔设计 +- [[00-索引]] — 返回文档总览 diff --git a/hzh/Gen2D/13-标签驱动提示词工程.md b/hzh/Gen2D/13-标签驱动提示词工程.md index 3d87157..483ca3b 100644 --- a/hzh/Gen2D/13-标签驱动提示词工程.md +++ b/hzh/Gen2D/13-标签驱动提示词工程.md @@ -1,10 +1,17 @@ -# 13. 标签驱动提示词工程 - -> **一句话概括**:40+ 预定义标签映射到精确的图像生成指令,保障风格一致性与管线友好性。 - +--- +tags: [prompt-engineering, tag-mapping, ai, llm, fallback-pattern, template] +create time: 2026-06-03 11:00 --- -## 处理流程 +# 13. 标签驱动提示词工程 + +## 概述 + +40+ 预定义标签映射到精确的图像生成指令,保障风格一致性与管线友好性。 + +## 正文 + +### 处理流程 ```mermaid flowchart LR @@ -24,11 +31,11 @@ flowchart LR --- -## 标签分类体系 +### 标签分类体系 系统内置 40+ 预定义标签,分为 **7 大类别**,每条标签精确映射到一条图像生成指令。 -### 内容类型 — 决定布局与格式 +#### 内容类型 — 决定布局与格式 | 标签 | 映射指令要点 | |------|-------------| @@ -40,7 +47,7 @@ flowchart LR | `序列帧` | 连续动画帧,网格排列,标注方向和帧数 | | `纸娃娃部件` | 可组合散件,统一比例和锚点 | -### 美术风格 — 决定渲染技术 +#### 美术风格 — 决定渲染技术 | 标签 | 映射指令要点 | |------|-------------| @@ -50,7 +57,7 @@ flowchart LR | `矢量` | 干净几何形状,平滑曲线 | | `扁平` | 无阴影或极少阴影,纯色块面 | -### 色调配色 +#### 色调配色 | 标签 | 映射指令要点 | |------|-------------| @@ -60,7 +67,7 @@ flowchart LR | `柔和` | 低饱和度,温和内敛 | | `单色` | 单一色相,明暗层次 | -### 线条粗细 +#### 线条粗细 | 标签 | 映射指令要点 | |------|-------------| @@ -69,7 +76,7 @@ flowchart LR | `中等` | 1-2px,清晰明确 | | `粗线` | 2-4px,粗犷有力 | -### 场景氛围 +#### 场景氛围 | 标签 | 映射指令要点 | |------|-------------| @@ -80,7 +87,7 @@ flowchart LR | `水下` | 珊瑚、水草、气泡 | | `沙漠` | 沙丘、仙人掌、绿洲 | -### 光照效果 +#### 光照效果 | 标签 | 映射指令要点 | |------|-------------| @@ -89,7 +96,7 @@ flowchart LR | `戏剧` | 强烈明暗对比,聚光灯效果 | | `霓虹` | 高饱和彩色光源,赛博朋克辉光 | -### 情绪基调 +#### 情绪基调 | 标签 | 映射指令要点 | |------|-------------| @@ -101,7 +108,7 @@ flowchart LR --- -## 精灵图特殊处理 +### 精灵图特殊处理 ```mermaid flowchart TD @@ -121,9 +128,13 @@ flowchart TD 这些指令对下游的 `SplitSprite` 切割算法至关重要。 +> [!question] 为什么精灵图需要特殊的网格布局指令? +> +> 如果不指定间隙要求,AI 生成的图片往往会让人物之间几乎没有空隙——人类艺术家这样做是为了最大化利用画布,但机器无法从中准确推断分割边界。**强制纯白间隙 = 人为制造"裂缝"**,让投影检测算法可以像翻书页一样逐页分离内容。这就是典型的「用约束换精度」的设计思路。 + --- -## 未知标签回退 +### 未知标签回退 当用户输入系统未预定义的标签时,不会报错,而是降级为通用风格描述: @@ -141,7 +152,7 @@ func tagToInstruction(tag string) string { --- -## PromptOptimizer 节点 +### PromptOptimizer 节点 PromptOptimizer 是 Eino 管线的第一个节点,负责将标签指令组装为最终提示词。 @@ -186,13 +197,13 @@ flowchart LR --- -## LLM 回退机制 +### LLM 回退机制 当 LLM 不可用时(API Key 未配置 / 网络故障),系统自动降级到模板生成: ```mermaid flowchart TD - A["callLLMRefine"] --> B{"API Key\n已配置?"} + A["callLLMRefine"] --> B{"API Key 已配置?"} B -->|"否"| F["fallbackRefine\n模板回退"] B -->|"是"| C["调用 Chat API"] C --> D{"调用成功?"} @@ -204,13 +215,9 @@ flowchart TD 模板回退同样遵循标签驱动逻辑,保证即使没有 LLM 参与,生成的提示词也具备结构化和一致性。 ---- - ## 关联文档 -| 文档 | 关系 | -|------|------| -| [05 - 生成管线](05-generation-pipeline.md) | PromptOptimizer 是管线第一阶段 | -| [06 - 精灵图处理](06-sprite-processing.md) | 网格布局指令影响切割算法 | -| [01 - 系统总览](01-system-overview.md) | 提示词工程在整体架构中的位置 | -| [00 - 索引](00-index.md) | 返回文档总览 | +- [[05-生成管线]] — PromptOptimizer 是管线第一阶段 +- [[06-精灵图处理]] — 网格布局指令影响切割算法 +- [[01-系统总览]] — 提示词工程在整体架构中的位置 +- [[00-索引]] — 返回文档总览 diff --git a/hzh/Gen2D/14-部署架构.md b/hzh/Gen2D/14-部署架构.md index 01d0f1b..318c09c 100644 --- a/hzh/Gen2D/14-部署架构.md +++ b/hzh/Gen2D/14-部署架构.md @@ -1,10 +1,17 @@ -# 14. 部署架构 - -> **一句话概括**:Docker Compose 编排 + Prometheus 监控 + Grafana 可视化 + 自动化部署脚本。 - +--- +tags: [deployment, docker-compose, monitoring, prometheus, grafana, devops] +create time: 2026-06-03 11:05 --- -## 容器拓扑 +# 14. 部署架构 + +## 概述 + +Docker Compose 编排 + Prometheus 监控 + Grafana 可视化 + 自动化部署脚本。 + +## 正文 + +### 容器拓扑 ```mermaid flowchart TB @@ -49,9 +56,7 @@ flowchart TB style GF fill:#f46800,stroke:#333,color:#fff ``` ---- - -## 服务组成 +### 服务组成 | 服务 | 镜像 | 端口 | 职责 | |------|------|------|------| @@ -64,9 +69,9 @@ flowchart TB --- -## 监控栈 +### 监控栈 -### Prometheus 配置 +#### Prometheus 配置 ```yaml # deploy/prometheus/prometheus.yml @@ -80,7 +85,7 @@ scrape_configs: Prometheus 每 **10 秒**抓取一次后端的 `/metrics` 端点,采集全部 35+ 指标。 -### Grafana 自动化 +#### Grafana 自动化 ```mermaid flowchart LR @@ -101,7 +106,7 @@ Grafana 通过 provisioning 机制自动加载: - **数据源配置** — 指向 Prometheus 实例 - **仪表盘 JSON** — 预定义的 3 个仪表盘 -### 告警规则 +#### 告警规则 10 条告警规则覆盖全栈关键指标: @@ -120,7 +125,7 @@ Grafana 通过 provisioning 机制自动加载: --- -## 部署脚本 +### 部署脚本 ```mermaid flowchart TD @@ -142,7 +147,7 @@ flowchart TD --- -## 配置管理 +### 配置管理 ```mermaid flowchart LR @@ -174,13 +179,13 @@ flowchart LR --- -## 网络与存储 +### 网络与存储 -### Docker 网络 +#### Docker 网络 所有服务加入 `gen2d-v2-net` 桥接网络,容器间通过服务名互相访问。 -### 数据卷 +#### 数据卷 | 卷名 | 挂载点 | 用途 | |------|--------|------| @@ -188,13 +193,9 @@ flowchart LR | MySQL data | 默认 | 数据库持久化 | | Redis data | 默认 | 缓存持久化 | ---- - ## 关联文档 -| 文档 | 关系 | -|------|------| -| [07 - 可观测性](07-observability.md) | Prometheus 指标与 Grafana 仪表盘详情 | -| [15 - 配置级联](15-config-cascade.md) | YAML / ENV / Default 三层配置机制 | -| [01 - 系统总览](01-system-overview.md) | 部署架构在整体系统中的位置 | -| [00 - 索引](00-index.md) | 返回文档总览 | +- [[07-可观测性]] — Prometheus 指标与 Grafana 仪表盘详情 +- [[15-配置级联机制]] — YAML / ENV / Default 三层配置机制 +- [[01-系统总览]] — 部署架构在整体系统中的位置 +- [[00-索引]] — 返回文档总览 diff --git a/hzh/Gen2D/15-配置级联机制.md b/hzh/Gen2D/15-配置级联机制.md index f8d92ce..ee8c5e9 100644 --- a/hzh/Gen2D/15-配置级联机制.md +++ b/hzh/Gen2D/15-配置级联机制.md @@ -1,10 +1,17 @@ -# 15. 配置级联机制 - -> **一句话概括**:Viper 三层配置级联 — YAML 文件 -> 环境变量 -> 默认值,一处配置随处运行。 - +--- +tags: [config, viper, environment-variables, yaml, deployment, configuration-management] +create time: 2026-06-03 11:10 --- -## 配置加载流程 +# 15. 配置级联机制 + +## 概述 + +Viper 三层配置级联 — YAML 文件 -> 环境变量 -> 默认值,一处配置随处运行。 + +## 正文 + +### 配置加载流程 ```mermaid flowchart LR @@ -41,7 +48,7 @@ flowchart LR --- -## 配置结构体 +### 配置结构体 ```go type Config struct { @@ -73,7 +80,7 @@ type Config struct { --- -## 环境变量绑定 +### 环境变量绑定 每个配置字段都有对应的环境变量绑定,命名规则为 `GEN2D_` 前缀 + 大写下划线格式: @@ -104,7 +111,7 @@ func bindEnvVars(v *viper.Viper) { --- -## 使用场景 +### 使用场景 ```mermaid flowchart TD @@ -140,12 +147,12 @@ flowchart TD --- -## 加载过程详解 +### 加载过程详解 ```mermaid sequenceDiagram participant Main as main.go - participant Viper + participant Viper as Viper participant YAML as config.yaml participant ENV as 环境变量 participant Cfg as Config Struct @@ -175,7 +182,7 @@ sequenceDiagram --- -## 与部署的关系 +### 与部署的关系 配置级联机制与部署架构紧密配合: @@ -185,13 +192,9 @@ sequenceDiagram | `docker compose` | `.env` 文件注入 | `env_file: ./backend/.env` | | Kubernetes | ConfigMap + Secret | `GEN2D_DSN` 从 Secret 注入 | ---- - ## 关联文档 -| 文档 | 关系 | -|------|------| -| [14 - 部署架构](14-deployment.md) | 配置管理在部署中的应用 | -| [01 - 系统总览](01-system-overview.md) | 配置在启动流程中的位置 | -| [02 - 协程池](02-worker-pool.md) | WorkerPoolConfig 控制并发参数 | -| [00 - 索引](00-index.md) | 返回文档总览 | +- [[14-部署架构]] — 配置管理在部署中的应用 +- [[01-系统总览]] — 配置在启动流程中的位置 +- [[02-协程池]] — WorkerPoolConfig 控制并发参数 +- [[00-索引]] — 返回文档总览