18 KiB
18 KiB
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):
{
"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 | 否 | 用户自定义提示词 |
响应:
{
"taskId": "550e8400-e29b-41d4-a716-446655440000"
}
处理流程:
- 查模型配置 — 根据
modelName查询model_gateway_models获取模型定义 - 创建构建记录 — 写入
model_gateway_build_record,状态为处理中 - 异步执行构建(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,存入构建记录
- 查找当前用户的对话模型(
- 推理模型(type=100):
- 回写结果 — 更新构建记录状态(成功/失败 + 耗时)
- 触发回调 — 若有
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):
{
"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 | 否 | 轮次 ID(prompts-core 场景) |
buildModelName |
string | 否 | 构建阶段使用的模型名(prompts-core 场景) |
响应:
{
"taskId": "550e8400-e29b-41d4-a716-446655440000"
}
处理流程:
- 鉴权与参数校验 — 获取用户信息,校验模型配置是否存在且已启用
- 队列门控(可选) — 通过 Redis Lua 脚本做严格的分布式队列插槽校验,避免并发超限
- 写入任务记录 — 状态为
1(执行中),记录requestPayload、callbackUrl、bizName等 - 操作日志 — 记录创建任务的审计日志到
model_gateway_logs_op - 计费预处理 — 若模型配置了
billing_config,从请求参数中提取计费维度数据 - 异步执行(goroutine
AsyncWorker.handleOne):
检查租户余额 ──> 调用模型服务 ──> 解析响应映射 ──> 计费计算并扣费
│
┌── 同步模式(call_mode=0):直接调用,等待 HTTP 响应
├── 异步模式(call_mode=1):提交任务后通过 QueryConfig 轮询结果
└── 流式模式(call_mode=2):通过 StreamConfig 解析 SSE 事件流
│
┌── prompts-core 场景:对模型输出做 ParseAndValidate + 重试机制
└── 普通场景:直接映射响应
│
上传 OSS ──> 更新任务为成功 ──> 触发回调
- 重试机制:超时、内部错误等可重试错误会按
retry_times重试;硬错误(参数错误等)直接失败 - 回调通知:任务成功或失败后,若有
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) |
| 数据库 | PostgreSQL(GoFrame pgsql 驱动) |
| 缓存/队列 | Redis(分布式信号量、队列门控、运行时调参存储) |
| 服务发现 | Consul |
| 链路追踪 | Jaeger(OTLP 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(可选,用于服务发现)
本地开发
- 克隆仓库
- 在 PostgreSQL 中执行
update.sql创建所有表 - 修改
config.yml中的数据库、Redis 和 Consul 配置 - 启动服务:
go run main.go - 服务默认监听
3004端口(可在config.yml中修改)
Docker 构建
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 | 已下载 | 批量查询接口标记 |
接口文档
地址:https://s.apifox.cn/42764602-cb82-45fa-887f-c751311ba406
License
内部项目 — 红未来科技