217 lines
7.6 KiB
Markdown
217 lines
7.6 KiB
Markdown
# Essence 报告:AI Pipeline 架构
|
||
|
||
**项目:** Lumina
|
||
**Lens(透镜):** Mechanical(技术实现)
|
||
**设计分析:** Article AI Pipeline Service
|
||
**文件数:** 6 个核心文件
|
||
**状态:** 完成
|
||
|
||
---
|
||
|
||
## 定位
|
||
|
||
** standout design:** 模块化 AI 内容处理管道,支持配置化 Prompt 和 Model。
|
||
|
||
**为什么值得学:**
|
||
- 解决「AI 处理步骤硬编码」的通病
|
||
- 结构化输出协议可复用到其他 AI 项目
|
||
- 任务状态机设计可迁移到任意异步处理系统
|
||
|
||
---
|
||
|
||
## 核心文件
|
||
|
||
| 文件 | 角色 | 关键内容 |
|
||
|------|------|----------|
|
||
| `article_ai_pipeline_service.py:102` | Pipeline 编排器 | `ArticleAIPipelineService` 主类 |
|
||
| `article_ai_pipeline_service.py:114` | 输出契约定义 | `STRUCTURED_OUTPUT_CONTRACTS` |
|
||
| `article_ai_pipeline_service.py:2039` | 清洗阶段 | `process_article_cleaning` |
|
||
| `article_ai_pipeline_service.py:3567` | AI 内容生成 | `process_ai_content` |
|
||
| `task_state.py:1` | 状态机 | `ALLOWED_TASK_STATUS_TRANSITIONS` |
|
||
|
||
---
|
||
|
||
## Call Chain 调用链
|
||
|
||
```
|
||
submit_article(url)
|
||
│
|
||
▼
|
||
┌─────────────────────────┐
|
||
│ ingest_service │
|
||
│ - fetch raw HTML │
|
||
└───────────┬─────────────┘
|
||
│
|
||
▼
|
||
┌─────────────────────────┐ ┌─────────────────────────┐
|
||
│ process_article_cleaning│────▶│ process_article_validate│
|
||
│ - 清洗原始内容为 Markdown│ │ - 验证内容合规性 │
|
||
└───────────┬─────────────┘ └───────────┬─────────────┘
|
||
│ │
|
||
│ ▼
|
||
│ ┌─────────────────────────┐
|
||
│ │ process_article_classify│
|
||
│ │ - 自动分类 │
|
||
│ └───────────┬─────────────┘
|
||
│ │
|
||
▼ ▼
|
||
┌─────────────────────────┐ ┌─────────────────────────┐
|
||
│ process_article_tagging │ │ 其他 AI 任务... │
|
||
│ - 自动打标签 │ │ │
|
||
└───────────┬─────────────┘ └─────────────────────────┘
|
||
│
|
||
▼
|
||
┌─────────────────────────┐
|
||
│ process_ai_content │
|
||
│ - summary/key_points/ │
|
||
│ quotes/outline/ │
|
||
│ infographic │
|
||
└─────────────────────────┘
|
||
```
|
||
|
||
---
|
||
|
||
## 核心设计模式
|
||
|
||
### 1. 结构化输出契约(Structured Output Contract)
|
||
|
||
**文件:** `article_ai_pipeline_service.py:114-205`
|
||
|
||
```python
|
||
STRUCTURED_OUTPUT_CONTRACTS = {
|
||
"classification": PromptOutputContract(
|
||
mode="structured_json",
|
||
response_format={
|
||
"type": "json_schema",
|
||
"json_schema": {
|
||
"name": "article_classification_result",
|
||
"schema": {"type": "object", "properties": {
|
||
"category_id": {"type": "string"}
|
||
}, "required": ["category_id"]}
|
||
}
|
||
},
|
||
system_instruction="固定输出协议:必须返回单个 JSON 对象..."
|
||
),
|
||
# ... tagging, validation, outline
|
||
}
|
||
```
|
||
|
||
**作用:** 将「AI 输出格式不稳定的」问题在系统设计层面解决,不再依赖 Prompt Engineering。
|
||
|
||
### 2. 可配置 Prompt + Model 绑定
|
||
|
||
**文件:** `article_ai_pipeline_service.py:255-334`
|
||
|
||
**查询优先级:**
|
||
1. 传入的 `model_config_id` / `prompt_config_id`
|
||
2. Category 级别的配置
|
||
3. 全局默认配置
|
||
|
||
**迁移价值:**
|
||
- 无需重启服务即可切换模型
|
||
- A/B 测试不同 Prompt 效果
|
||
- 支持多模型混用(强项模型处理特定任务)
|
||
|
||
### 3. Pipeline 状态机
|
||
|
||
**文件:** `task_state.py:12-23`
|
||
|
||
```python
|
||
ALLOWED_TASK_STATUS_TRANSITIONS = {
|
||
TASK_STATUS_PENDING: {TASK_STATUS_PROCESSING, TASK_STATUS_CANCELLED},
|
||
TASK_STATUS_PROCESSING: {
|
||
TASK_STATUS_PENDING, TASK_STATUS_COMPLETED, TASK_STATUS_FAILED
|
||
},
|
||
TASK_STATUS_FAILED: {TASK_STATUS_PENDING}, # 支持重试
|
||
TASK_STATUS_CANCELLED: {TASK_STATUS_PENDING}, # 取消后可恢复
|
||
}
|
||
```
|
||
|
||
**扩展点:** 状态流转规则集中定义,新增状态只需改一行。
|
||
|
||
### 4. 续写/修复机制
|
||
|
||
**文件:** `article_ai_pipeline_service.py:3567-3598`
|
||
|
||
```python
|
||
async def process_ai_content(
|
||
self,
|
||
article_id: str,
|
||
content_type: str,
|
||
continuation_source_usage_id: str | None = None, # 续写源头
|
||
continuation_feedback: str | None = None, # 用户反馈
|
||
...
|
||
):
|
||
```
|
||
|
||
**设计意图:** 不是一轮生成完事,而是支持「反馈 → 修复 → 续写」的迭代循环。
|
||
|
||
---
|
||
|
||
## Pattern 分析
|
||
|
||
| 维度 | 内容 |
|
||
|------|------|
|
||
| **问题** | 如何为同一篇文章批量生成摘要、标签、分类、翻译等 AI 内容? |
|
||
| **模式** | 配置驱动的 Pipeline + 结构化输出契约 |
|
||
| **替代方案** | A) 每个功能独立 endpoint,各自调用 AI(重复配置);B) 固定流程硬编码(不可配置) |
|
||
| **Tradeoff** | 放弃实时性(Pipeline 异步执行),换取可配置性和失败隔离 |
|
||
| **证据** | `article_ai_pipeline_service.py:102-205` 定义契约;`task_state.py:1-60` 状态机 |
|
||
|
||
---
|
||
|
||
## Migration 示例(可执行)
|
||
|
||
**场景:** 为自己的项目实现「User Content → AI Processed Content」管道。
|
||
|
||
```python
|
||
# pipeline.py - 核心骨架(≤20 行)
|
||
from dataclasses import dataclass
|
||
from typing import Callable
|
||
|
||
@dataclass
|
||
class PipelineStep:
|
||
name: str
|
||
processor: Callable[[str], str]
|
||
output_type: str = "text"
|
||
|
||
class ContentPipeline:
|
||
def __init__(self):
|
||
self.steps: list[PipelineStep] = []
|
||
|
||
def add_step(self, step: PipelineStep):
|
||
self.steps.append(step)
|
||
|
||
async def process(self, content: str) -> dict:
|
||
results = {"input": content}
|
||
for step in self.steps:
|
||
results[step.name] = await step.processor(results.get(step.output_type, content))
|
||
return results
|
||
|
||
# 使用
|
||
pipeline = ContentPipeline()
|
||
pipeline.add_step(PipelineStep("clean", clean_html, "text"))
|
||
pipeline.add_step(PipelineStep("summarize", generate_summary, "clean"))
|
||
```
|
||
|
||
---
|
||
|
||
## Pitfalls 陷阱
|
||
|
||
| 陷阱 | 原因 | 规避方法 |
|
||
|------|------|----------|
|
||
| AI 输出格式不稳定 | 未定义结构化契约 | 复制 `STRUCTURED_OUTPUT_CONTRACTS` 思想 |
|
||
| Pipeline 某步失败导致整体失败 | 无状态隔离 | 每步独立状态,失败可重试 |
|
||
| Prompt 调优困难 | Prompt 与代码耦合 | 数据库存储,支持按 category 配置 |
|
||
| 成本不可控 | 无 Token 统计 | 每步记录 `price_input_per_1k` / `price_output_per_1k` |
|
||
|
||
---
|
||
|
||
## 证据核对
|
||
|
||
- [x] 设计真实存在(`article_ai_pipeline_service.py:102`)
|
||
- [x] 文件列表 ≤ 10(实际 6 个)
|
||
- [x] Pattern 可解释(配置驱动 + 结构化契约)
|
||
- [x] Migration ≤ 20 行(上面示例 19 行)
|
||
- [x] Pitfalls 具体(有代码/配置对应)
|