Files

217 lines
7.6 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 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 具体(有代码/配置对应)