Files
model-gateway/README.md
T
WangLiZhao 6a48edb742 docs(readme): 更新项目文档为智能模型网关
chore(database): 优化模型配置表结构和索引设计

- 重新设计 model_gateway_models 表的字段结构和约束
- 新增 billing_config 和 special_params 配置字段
- 优化索引策略,包括复合索引和条件索引的创建
- 统一表注释格式和字段说明文档
- 移除过时的配置字段和冗余索引
2026-07-07 09:47:30 +08:00

435 lines
18 KiB
Markdown
Raw 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.
# Model Gateway — 智能模型网关
统一的 AI 模型网关服务,负责多模态(文本、图像、音频、向量、视频、全模态)模型调用的路由、执行与管理。基于 Go (GoFrame v2) 构建,使用 PostgreSQL、Redis 和 Consul 服务发现。
---
## 架构概览
```
客户端 ──> HTTP API ──> Controller ──> Service ──> DAO ──> PostgreSQL
│ └─> Redis(信号量/队列)
├──> AI 模型服务商(OpenAI、阿里云、火山引擎等)
├──> OSS 文件服务
├──> Admin-go(租户/余额)
├──> Prompts-core(会话回调)
└──> Skill 技能服务
```
系统分为六个层次:
| 层次 | 目录 | 职责 |
|---|---|---|
| **控制器层** | `controller/` | HTTP 路由注册与请求绑定 |
| **业务逻辑层** | `service/` | 模型配置 CRUD、异步任务执行、提示词构建、队列管理、统计 |
| **数据访问层** | `dao/` | PostgreSQL 数据访问(GoFrame ORM |
| **模型层** | `model/` | DTO 请求/响应结构与数据库实体定义 |
| **通用工具层** | `common/` | 计费、类型转换、文件处理、请求头、映射、提示词、流式解析等工具 |
| **常量层** | `consts/` | 公共常量和表名定义 |
---
## 核心功能
### 1. 模型配置管理 (`model_gateway_models`)
提供 AI 模型服务配置的完整增删改查:
| 字段 | 说明 |
|---|---|
| `model_name` | 模型名称(唯一标识) |
| `model_type` | 模型分类(100=推理, 200=图像, 300=音频, 400=向量, 500=全模态, 600=视频) |
| `operator_name` | 运营商标识(OpenAI、阿里云、火山引擎等) |
| `base_url` | 模型服务地址 |
| `http_method` | GET / POST |
| `head_msg` | 每次调用注入的请求头 |
| `form_json` | 动态表单定义(用于前端按模型渲染参数表单) |
| `request_mapping` / `response_mapping` | 标准格式与提供商 API 之间的参数映射 |
| `call_mode` | 0=同步, 1=异步, 2=流式 |
| `max_concurrency` | 单模型最大并发数(按租户维度限流) |
| `timeout_seconds` | 请求超时时间(秒) |
| `retry_times` | 失败重试次数 |
| `billing_config` | 计费规则(推理阶梯计价 / 视频分辨率计价) |
| `stream_config` | SSE 流式解析配置 |
| `query_config` | 异步任务轮询/查询配置 |
**支持的模型类型:**
| 类型码 | 类别 | 子类型 |
|---|---|---|
| 100 | 推理模型 | 文本生成、对话 |
| 200-205 | 图像模型 | 文生图、图生图、图片编辑、图片变体、图文生图 |
| 300-303 | 音频模型 | 文生语音、语音转文字、语音转语音 |
| 400-402 | 向量模型 | 文本嵌入、重排序 |
| 500-502 | 全模态模型 | 文图音频、视觉理解 |
| 600-604 | 视频模型 | 文生视频、图生视频、图文生视频、视频生视频 |
### 2. 异步任务执行 (`model_gateway_task`)
任务生命周期如下:
```
CreateTask ──> state=0(排队中)
Worker 抢占(state=1,执行中)
├── 同步模式:直接调用模型
├── 异步模式:提交任务,通过 QueryConfig 轮询
└── 流式模式:通过 StreamConfig 解析 SSE 事件
├── 成功(state=2)──> 上传结果到 OSS
│ └──> 触发回调(如有配置)
└── 失败(state=3)──> 重试(最多 retry_times 次)
客户端下载(state=4,已下载)
```
**关键特性:**
- **任务创建**:立即返回 `taskId`,执行在 goroutine 中异步完成
- **并发控制**:基于 Redis 的分布式信号量,按模型维度限制并发数
- **队列门控**:基于 Redis Lua 脚本的严格队列插槽机制,防止分布式创建下超限
- **自动重试**:临时性错误(超时、内部错误)自动重试;硬错误直接失败
- **OSS 上传**:结果以 JSON 格式上传到 OSS 文件服务,OSS URL 存入任务记录
- **回调通知**:任务完成时可选触发 HTTP 回调(`TriggerCallback``TriggerPromptsCallback``CallbackBuildResult`
- **阶梯计费**:支持推理 token 阶梯计价和视频分辨率计价两种模型
### 3. 提示词构建 (`/buildMessages`)
异步构建推理模型的结构化消息:
- 从配置合并系统提示词(`modelPrompts.types`
- 拉取技能 Markdown 内容(`SkillMdContent`
- 通过 `GetSessionHistory` 注入会话历史
- 根据附件数量自动拆分多轮(`SplitByAttachment`
- 支持视频模型通过对话模型编排的多轮生成
### 4. 动态调参 (`/autoTune`)
周期性自动调参(建议每小时触发),基于近期 P90 执行耗时和到达率动态调整 `max_concurrency` 和队列上限:
- 读取模型配置(数据库中的上限值 cap)
- 根据时间窗口内已完成任务统计 P90 执行耗时
- 使用 Little 定律(到达率 × P90 / 利用率)计算新的并发数
- 单次调整幅度限制在 ±50%
- 运行时参数写入 Redis(2 小时 TTL),不修改数据库中的上限值
### 5. 日统计 (`model_gateway_log_stat`)
按天、租户、创建人、模型维度的请求计数,使用原子 upsert(`ON DUPLICATE KEY UPDATE`)。
按天、租户、创建人、模型维度的请求计数,使用原子 upsert(`ON DUPLICATE KEY UPDATE`)。
### 6. 核心接口详细说明
#### 6.1 构建消息结构 `/buildMessages`
构建推理模型的提示词消息体,将系统提示词、技能知识、会话历史、用户自定义提示词合并为最终的消息结构。
**请求示例(JSON):**
```json
{
"modelName": "gpt-4o",
"buildType": 1,
"skillName": "code-review",
"callbackUrl": "http://callback.example.com/result",
"nodeId": "node_001",
"sessionId": "session_abc",
"customPrompt": "请用中文回答",
"messages": {
"model": "gpt-4o",
"max_tokens": 4096,
"messages": [
{"role": "system", "content": [{"type": "text", "text": "你是一名助手"}]},
{"role": "user", "content": [{"type": "text", "text": "写一封邮件"}]}
]
}
}
```
| 请求字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
| `modelName` | string | 是 | 网关模型名称 |
| `buildType` | int | 是 | 构建类型:1=单轮构建, 2=多轮构建 |
| `skillName` | string | 否 | 技能名称,用于拉取技能 MD 知识 |
| `callbackUrl` | string | 否 | 构建完成后的回调地址 |
| `nodeId` | string | 否 | 节点 ID(用于查询会话历史) |
| `sessionId` | string | 否 | 会话 ID |
| `messages` | object | 是 | 前端构建的消息结构 |
| `customPrompt` | string | 否 | 用户自定义提示词 |
**响应:**
```json
{
"taskId": "550e8400-e29b-41d4-a716-446655440000"
}
```
**处理流程:**
1. **查模型配置** — 根据 `modelName` 查询 `model_gateway_models` 获取模型定义
2. **创建构建记录** — 写入 `model_gateway_build_record`,状态为处理中
3. **异步执行构建**goroutine):
- **推理模型(type=100**
- 从配置读取系统提示词(`modelPrompts.types.100`
- 若有 `skillName`,通过 `GetSkillUser` 拉取技能 ZIP,提取 MD 内容
- 合并系统提示词 + 技能知识 + 用户自定义提示词到 `messages`
- 通过 `GetSessionHistory` 注入会话历史
- 按附件数量自动拆分多轮(`SplitByAttachment`,按 video/image/audio 的 `maxCount` 约束)
- 返回单轮或多轮消息结构
- **视频模型(type=600-699**
- 查找当前用户的对话模型(`is_chat_model=1`
- 获取该对话模型的协议模板(`prompts_provider_protocol`
- 构建对话模型请求体,调用对话模型生成视频分镜的 JSON rounds
- 将 rounds 上传 OSS,存入构建记录
4. **回写结果** — 更新构建记录状态(成功/失败 + 耗时)
5. **触发回调** — 若有 `callbackUrl`POST 回调通知任务完成
**对应源码:**
- Controller[controller/model_gateway_task_controller.go]`ModelGatewayTask.BuildMessages`
- Service[service/task/task_service.go]`taskService.BuildMessages``executeBuild``buildResult`
- Util(提示词合并):[common/util/prompt.go]`MergePrompt``InjectHistory``SplitByAttachment`
- Util(技能知识):[service/prompt/prompt_files_handle_service.go]`SkillMdContent`
#### 6.2 创建异步任务 `/createTask`
创建异步模型调用任务,立即返回 `taskId`,后端 goroutine 异步执行模型调用、结果上传与计费。
**请求示例(JSON):**
```json
{
"modelName": "dall-e-3",
"bizName": "image-generator",
"callbackUrl": "http://callback.example.com/result",
"epicycleId": 12345,
"buildModelName": "gpt-4o",
"requestPayload": {
"model": "dall-e-3",
"prompt": "一只站在树枝上的猫头鹰,水墨风格",
"n": 1,
"size": "1024x1024"
}
}
```
| 请求字段 | 类型 | 必填 | 说明 |
|---|---|---|---|
| `modelName` | string | 是 | 模型名称 |
| `bizName` | string | 否 | 业务名称,用于统计区分 |
| `callbackUrl` | string | 否 | 任务完成后的回调地址 |
| `requestPayload` | object | 是 | 透传给模型服务的请求参数 |
| `epicycleId` | int64 | 否 | 轮次 IDprompts-core 场景) |
| `buildModelName` | string | 否 | 构建阶段使用的模型名(prompts-core 场景) |
**响应:**
```json
{
"taskId": "550e8400-e29b-41d4-a716-446655440000"
}
```
**处理流程:**
1. **鉴权与参数校验** — 获取用户信息,校验模型配置是否存在且已启用
2. **队列门控(可选)** — 通过 Redis Lua 脚本做严格的分布式队列插槽校验,避免并发超限
3. **写入任务记录** — 状态为 `1(执行中)`,记录 `requestPayload``callbackUrl``bizName`
4. **操作日志** — 记录创建任务的审计日志到 `model_gateway_logs_op`
5. **计费预处理** — 若模型配置了 `billing_config`,从请求参数中提取计费维度数据
6. **异步执行**goroutine `AsyncWorker.handleOne`):
```
检查租户余额 ──> 调用模型服务 ──> 解析响应映射 ──> 计费计算并扣费
┌── 同步模式(call_mode=0):直接调用,等待 HTTP 响应
├── 异步模式(call_mode=1):提交任务后通过 QueryConfig 轮询结果
└── 流式模式(call_mode=2):通过 StreamConfig 解析 SSE 事件流
┌── prompts-core 场景:对模型输出做 ParseAndValidate + 重试机制
└── 普通场景:直接映射响应
上传 OSS ──> 更新任务为成功 ──> 触发回调
```
7. **重试机制**:超时、内部错误等可重试错误会按 `retry_times` 重试;硬错误(参数错误等)直接失败
8. **回调通知**:任务成功或失败后,若有 `callbackUrl` 则 POST 回调;prompts-core 场景额外触发 `TriggerPromptsCallback`
**调用模式对比:**
| 模式 | 值 | 行为 |
|---|---|---|
| 同步 | 0 | 直接 HTTP 调用模型,等待完整响应后返回 |
| 异步 | 1 | 提交任务后通过 `QueryConfig` 配置的轮询地址定时查询结果 |
| 流式 | 2 | 对 SSE 事件流按 `StreamConfig` 规则解析(支持 concat/base64_concat/collect/final |
**对应源码:**
- Controller[controller/model_gateway_task_controller.go]`ModelGatewayTask.CreateTask`
- Service[service/task/task_service.go]`taskService.Create`
- Worker[service/task/worker.go]`asyncWorker.handleOne``InvokeModel`
- 队列门控:[service/queue/queue_gate.go]`AcquireQueueSlot``ReleaseQueueSlot`
- 运行时调参:[service/queue/runtime_tune.go]`GetRuntimeMaxConcurrency``GetRuntimeQueueLimit`
- 异步轮询:[common/util/pull_task.go]`PullTaskResult`
- 流式解析:[common/util/streaming.go]`ParseStreamResponse`
- 计费:[common/util/billing.go]`CalculateBilling``ExtractRequestBilling`
### 7. 全量 API 接口一览
| 分类 | 接口路径 | 方法 | 说明 |
|---|---|---|---|
| **模型配置** | `/createModel` | POST | 创建模型配置 |
| | `/updateModel` | PUT | 更新模型配置 |
| | `/deleteModel` | DELETE | 删除模型配置 |
| | `/getModel` | GET | 获取模型详情 |
| | `/listModel` | GET | 模型列表(分页+筛选) |
| | `/listType` | GET | 模型类型列表 |
| | `/listOperator` | GET | 运营商列表 |
| | `/updateChatModel` | POST | 设置当前用户的对话模型 |
| | `/getIsChatModel` | GET | 获取当前对话模型 |
| | `/autoTune` | POST | 触发动态调参 |
| **任务管理** | `/createTask` | POST | 创建异步任务 |
| | `/jobTask` | POST | 定时批量任务处理器 |
| | `/getTaskResult` | GET | 获取单条任务结果 |
| | `/getTaskBatch` | POST | 批量查询任务(成功任务自动标记为已下载) |
| | `/listTask` | GET | 任务列表分页查询 |
| | `/modelCallback` | POST | 接收异步模型回调通知 |
| | `/queryPending` | GET | 轮询进行中的异步任务 |
| **提示词** | `/buildMessages` | POST | 构建结构化提示词消息(异步) |
| **统计** | `/listModelStat` | GET | 日使用量统计 |
--
## 技术栈
| 组件 | 技术选型 |
|---|---|
| 语言 | Go 1.26 |
| 框架 | GoFrame v2`github.com/gogf/gf/v2` |
| 数据库 | PostgreSQLGoFrame pgsql 驱动) |
| 缓存/队列 | Redis(分布式信号量、队列门控、运行时调参存储) |
| 服务发现 | Consul |
| 链路追踪 | JaegerOTLP HTTP |
| 模型调用 | 原生 HTTP(可定制请求头、鉴权、超时) |
| 文件存储 | OSS 文件上传服务(multipart 表单) |
## 数据库表
| 表名 | 说明 |
|---|---|
| `model_gateway_models` | 模型服务配置(动态表单、请求/响应映射、计费规则) |
| `model_gateway_task` | 异步任务记录(状态机、重试、OSS 结果文件) |
| `model_gateway_build_record` | 提示词构建记录 |
| `model_gateway_logs_op` | 操作审计日志 |
| `model_gateway_logs_stat` | 日维度请求量统计 |
| `prompts_provider_protocol` | 视频模型编排的提供商协议模板 |
## 配置说明
详见 `config.yml`。主要配置块:
- `database` — PostgreSQL 连接(双数据源:`default` + `model_gateway`
- `redis` — Redis 连接
- `consul` — Consul 地址
- `jaeger` — Jaeger OTLP HTTP 端点
- `queryPending` — 异步任务自动轮询配置(调试开关)
- `jobTask` — 批量处理器间隔、批大小、协程池大小
- `modelPrompts.types` — 各模型类型的系统提示词(100/200/300/400/500
- `nodePrompts` — 节点路由提示词模板
## 快速开始
### 环境要求
- Go 1.26+
- PostgreSQL
- Redis
- Consul(可选,用于服务发现)
### 本地开发
1. 克隆仓库
2. 在 PostgreSQL 中执行 `update.sql` 创建所有表
3. 修改 `config.yml` 中的数据库、Redis 和 Consul 配置
4. 启动服务:
```bash
go run main.go
```
5. 服务默认监听 `3004` 端口(可在 `config.yml` 中修改)
### Docker 构建
```bash
docker build -t model-gateway .
docker run -p 3004:3004 model-gateway
```
Dockerfile 使用多阶段构建,基于 `golang:alpine` 镜像,配置了国内 Go 代理,输出剥离调试信息的精简二进制。
---
## 项目目录结构
```
.
├── main.go # 入口:路由注册、自动执行器启动
├── config.yml # 应用配置
├── Dockerfile # 多阶段 Docker 构建
├── go.mod / go.sum # Go 模块依赖
├── update.sql # 数据库 DDL
├── controller/ # HTTP 接口控制器
│ ├── model_gateway_models_controller.go
│ ├── model_gateway_task_controller.go
│ └── model_gateway_logs_stat_controller.go
├── service/ # 业务逻辑层
│ ├── gateway/ # OSS 上传、回调、余额、技能/会话查询
│ ├── model/ # 模型 CRUD、对话模型管理
│ ├── prompt/ # 文件拉取、技能 Markdown、提示词构建
│ ├── queue/ # 动态调参、信号量、队列门控、运行时调参
│ ├── stat/ # 使用量统计
│ └── task/ # 任务创建、Worker 执行、异步结果处理
├── dao/ # 数据访问层
│ ├── model_gateway_models_dao.go
│ ├── model_gateway_task_dao.go
│ ├── model_gateway_build_record_dao.go
│ ├── model_gateway_logs_stat_dao.go
│ ├── model_gateway_logs_op_dao.go
│ └── provider_protocol_dao.go
├── model/ # 数据模型
│ ├── dto/ # 请求/响应结构体
│ └── entity/ # 数据库实体
├── common/ # 通用工具
│ └── util/ # 计费、类型转换、文件处理、请求头、映射、提示词、流式、异步轮询
└── consts/ # 常量定义
└── public/ # 模型类型、任务状态、运营商列表、表名
```
---
## 任务状态机
| 状态码 | 含义 | 转换来源 |
|---|---:|---|
| 0 | 排队中 | 创建任务或重试入队 |
| 1 | 执行中 | Worker 抢占成功 |
| 2 | 成功(已上传 OSS | Worker 执行完成 |
| 3 | 失败 | Worker 执行出错或超时 |
| 4 | 已下载 | 批量查询接口标记 |
---
## License
内部项目 — 红未来科技