RAGFlow 多租户扩展服务
技术方案

独立代理服务架构 — 零侵入 RAGFlow 源码 · 代理限流 · 扩展管理
V2 — 代理隔离架构(替代 V1 源码侵入方案)
📅 编制日期:2026-06-08 📦 基线版本:RAGFlow v0.25.6 🔒 密级:内部

📑 目录

  1. 方案演进:V1 → V2 架构对比
  2. V2 核心设计理念
  3. 总体架构设计
  4. 扩展服务技术选型
  5. 代理层设计:限流与请求路由
  6. 管理服务设计:法人/团队/权限
  7. 数据模型设计
  8. 向量数据库租户分区
  9. 空间模型与访问隔离
  10. API 设计
  11. 部署架构
  12. RAGFlow 升级兼容策略
  13. 数据迁移方案
  14. 实施路线图
  15. 关键风险与应对

1 方案演进:V1 → V2 架构对比

🔴 V1 方案:源码侵入式

  • 直接修改 RAGFlow 源码
  • 在 RAGFlow 内部添加中间件
  • 修改 RAGFlow 数据库表结构
  • 修改 RAGFlow API 路由
  • 与 RAGFlow 版本强耦合
  • 升级 RAGFlow 需大量合并冲突

🟢 V2 方案:独立代理服务

  • 独立 Python 服务,零侵入 RAGFlow
  • 代理层拦截请求实现限流/鉴权
  • 扩展服务自有数据库,不修改 RAGFlow 表
  • 新增管理 API,原始 API 透传代理
  • 与 RAGFlow 版本解耦
  • 升级 RAGFlow 仅需验证 API 兼容性
V2 核心思路:基于 RAGFlow 现有租户(Tenant)隔离能力,在其前面部署一个独立的 Python 扩展服务。该服务承担两大职责:
代理层:拦截 RAGFlow 原始 API 请求,实现限流、鉴权、租户上下文注入后转发给 RAGFlow;
管理服务:提供法人/用户组/团队/角色/空间等扩展管理功能,最终通过 RAGFlow 原生租户机制实现数据隔离。

2 V2 核心设计理念

原则描述实现方式
零侵入 不修改 RAGFlow 任何源码 独立服务 + 反向代理
代理限流 限速通过代理原始接口实现 代理层拦截解析/检索等主要接口
租户复用 最终访问基于原始租户隔离 扩展服务管理法人/团队,映射到 RAGFlow Tenant
升级友好 RAGFlow 升级无需合并代码 仅依赖 RAGFlow 公开 API 契约
独立演进 扩展服务可独立迭代 独立代码仓库、独立数据库、独立部署

2.1 RAGFlow 现有租户隔离能力复用

RAGFlow v0.25.6 已具备的租户隔离能力:

能力RAGFlow 原生实现V2 复用方式
租户数据隔离所有核心表通过 tenant_id FK 隔离直接复用,代理层注入 tenant_id
用户-租户关联UserTenant 表 + role 字段复用,扩展服务同步管理映射关系
知识库权限permission 字段(me/team)复用 team 模式,扩展服务控制成员可见性
API Token 认证Bearer Token + JWT代理层透传 Token,附加鉴权逻辑
团队邀请owner 可邀请 member扩展服务调用 RAGFlow 邀请 API

3 总体架构设计

3.1 系统架构总览

┌──────────────┐ │ 客户端/SDK │ └──────┬───────┘ │ 所有请求 ▼ ┌─────────────────────────────────────────────────────────────────────┐ │ RAGFlow Extension Gateway │ │ (独立 Python 服务 :9380) │ │ │ │ ┌────────────────────────────────────────────────────────────────┐ │ │ │ Nginx / Traefik 入口 │ │ │ └──────────────────────────┬─────────────────────────────────────┘ │ │ │ │ │ ┌───────────────┼───────────────┐ │ │ ▼ ▼ │ │ ┌─────────────────────┐ ┌─────────────────────────┐ │ │ │ 代理层 (Proxy) │ │ 管理服务 (Management) │ │ │ │ │ │ │ │ │ │ • 请求拦截/鉴权 │ │ • 法人管理 │ │ │ │ • 限流(滑动窗口) │ │ • 用户组管理 │ │ │ │ • 租户上下文注入 │ │ • 团队管理 │ │ │ │ • 请求转发/改写 │ │ • 角色/权限管理 │ │ │ │ • 响应过滤 │ │ • 空间管理 │ │ │ │ • 审计日志 │ │ • 限流规则管理 │ │ │ └────────┬────────────┘ │ • 文档管理(跨租户) │ │ │ │ │ • 用户授权/限流配置 │ │ │ │ 透传/改写后请求 └────────────┬────────────┘ │ │ ▼ │ │ │ ┌──────────────────┐ │ 调用 RAGFlow API │ │ │ RAGFlow 原始 │◄───────────────────────┘ │ │ │ API 服务 │ │ │ │ (:9380 原端口) │ │ │ └────────┬─────────┘ │ │ │ │ └───────────┼───────────────────────────────────────────────────────────┘ │ ▼ ┌───────────────────────────────────────────────────────────────────┐ │ RAGFlow 原始服务(不修改) │ │ │ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ │ │ ragflow- │ │ ES / │ │ MinIO │ │ Redis │ │ │ │ server │ │ Infinity │ │ │ │ │ │ │ └──────────┘ └──────────┘ └──────────┘ └──────────┘ │ │ ┌──────────┐ │ │ │ MySQL │ ← RAGFlow 原始数据库(不修改表结构) │ │ └──────────┘ │ └───────────────────────────────────────────────────────────────────┘ ┌───────────────────────────────────────────────────────────────────┐ │ 扩展服务自有数据库 │ │ ┌──────────┐ │ │ │ MySQL/ │ ← organization, user_group, role, permission, │ │ │ PG │ space, rate_limit_quota, audit_log ... │ │ └──────────┘ │ │ ┌──────────┐ │ │ │ Redis │ ← 限流计数器、权限缓存、会话状态 │ │ └──────────┘ │ └───────────────────────────────────────────────────────────────────┘

3.2 请求流转路径

客户端请求 Extension Gateway 路由判断 管理 API → 扩展服务处理
代理 API → 鉴权+限流 → 转发 RAGFlow

3.3 关键映射关系

核心映射:扩展服务的法人/用户组/团队层级,最终映射到 RAGFlow 的 Tenant。

Organization (法人) → 1:N → UserGroup (用户组) → 1:N → Team (团队)RAGFlow Tenant

每个团队对应一个 RAGFlow Tenant,扩展服务维护 team ↔ tenant_id 的映射表。
数据隔离完全依赖 RAGFlow 原生的 tenant_id 过滤机制,扩展服务不触碰 RAGFlow 数据。

4 扩展服务技术选型

组件选型理由
Web 框架FastAPI高性能异步、自动 OpenAPI 文档、与 RAGFlow (Quart) 无耦合
HTTP 代理httpx (async)异步 HTTP 客户端,支持连接池、流式转发
数据库 ORMSQLAlchemy 2.0 + async异步 ORM,独立数据库,不与 RAGFlow 共享
数据库PostgreSQL支持 RLS 行级安全,JSONB,独立于 RAGFlow 的 MySQL
缓存/限流Redis滑动窗口限流、权限缓存、可复用 RAGFlow 的 Redis 实例(不同 DB)
认证PyJWT + OAuth2验证 RAGFlow JWT + 扩展服务自有 JWT
部署Docker Compose与 RAGFlow 同一 compose 网络,独立容器
向量库分区ES / Milvus / Infinity复用 RAGFlow 已有向量库,通过 API 管理分区

4.1 项目结构

ragflow-extension/
├── app/
│   ├── main.py                    # FastAPI 入口
│   ├── config.py                  # 配置管理
│   ├── proxy/                     # 代理层
│   │   ├── router.py              # 代理路由(catch-all)
│   │   ├── middleware.py          # 鉴权/限流/审计中间件
│   │   ├── rate_limiter.py        # 滑动窗口限流器
│   │   ├── tenant_resolver.py     # 租户上下文解析
│   │   └── request_transformer.py # 请求改写(注入 tenant_id 等)
│   ├── management/                # 管理服务
│   │   ├── org_router.py          # 法人管理 API
│   │   ├── group_router.py        # 用户组管理 API
│   │   ├── team_router.py         # 团队管理 API
│   │   ├── role_router.py         # 角色/权限管理 API
│   │   ├── space_router.py        # 空间管理 API
│   │   ├── ratelimit_router.py    # 限流规则管理 API
│   │   └── doc_admin_router.py    # 系统管理员文档管理 API
│   ├── models/                    # 数据模型
│   │   ├── organization.py
│   │   ├── user_group.py
│   │   ├── team.py
│   │   ├── role.py
│   │   ├── permission.py
│   │   ├── space.py
│   │   ├── rate_limit_quota.py
│   │   └── audit_log.py
│   ├── services/                  # 业务逻辑
│   │   ├── org_service.py
│   │   ├── team_service.py
│   │   ├── permission_engine.py
│   │   ├── space_service.py
│   │   └── ragflow_client.py      # RAGFlow API 客户端
│   ├── db/                        # 数据库
│   │   ├── session.py
│   │   └── migrations/
│   └── schemas/                   # Pydantic 模型
├── docker/
│   ├── Dockerfile
│   └── docker-compose.yml         # 扩展服务 compose
├── tests/
├── pyproject.toml
└── README.mdProject Structure

5 代理层设计:限流与请求路由

5.1 代理层核心职责

职责实现方式拦截接口
限流Redis 滑动窗口计数器所有 API
鉴权JWT 校验 + 扩展权限引擎所有 API
租户路由X-Tenant-ID → tenant_id 映射所有 API
请求改写注入/替换 tenant_id、过滤字段数据写入 API
响应过滤按权限过滤返回数据数据读取 API
审计日志异步写入审计表敏感操作 API

5.2 代理路由实现

# app/proxy/router.py

from fastapi import APIRouter, Request, Depends
from fastapi.responses import StreamingResponse
import httpx

RAGFLOW_BASE = "http://ragflow-server:9380"

proxy_router = APIRouter()

MANAGEMENT_PREFIXES = [
    "/api/v1/organizations",
    "/api/v1/teams",
    "/api/v1/rate-limits",
    "/api/v1/admin/",
    "/api/v1/ext/",
]

@proxy_router.api_route("/{path:path}", methods=["GET", "POST", "PUT", "DELETE", "PATCH"])
async def proxy_to_ragflow(request: Request, path: str):
    full_path = f"/{path}"

    if any(full_path.startswith(p) for p in MANAGEMENT_PREFIXES):
        return Response(status_code=404)

    tenant_id = getattr(request.state, "tenant_id", None)
    user = getattr(request.state, "user", None)

    headers = dict(request.headers)
    headers.pop("host", None)
    headers.pop("content-length", None)

    if tenant_id:
        headers["X-Tenant-ID"] = tenant_id

    body = await request.body()

    url = f"{RAGFLOW_BASE}{full_path}"
    if request.url.query:
        url += f"?{request.url.query}"

    async with httpx.AsyncClient(timeout=300.0) as client:
        if request.method == "GET":
            resp = await client.get(url, headers=headers)
        elif request.method == "POST":
            resp = await client.post(url, headers=headers, content=body)
        elif request.method == "PUT":
            resp = await client.put(url, headers=headers, content=body)
        elif request.method == "DELETE":
            resp = await client.delete(url, headers=headers)
        elif request.method == "PATCH":
            resp = await client.patch(url, headers=headers, content=body)
        else:
            return Response(status_code=405)

    return Response(
        content=resp.content,
        status_code=resp.status_code,
        headers=dict(resp.headers),
    )Python — Proxy Router

5.3 限流中间件实现

# app/proxy/rate_limiter.py

import time
from typing import Optional

class SlidingWindowRateLimiter:
    def __init__(self, redis_client):
        self.redis = redis_client

    async def check(self, scope_type: str, scope_id: str,
                    api_path: str, max_requests: int, window: int = 60) -> bool:
        key = f"ratelimit:{scope_type}:{scope_id}:{api_path}"
        now = time.time()
        pipe = self.redis.pipeline()
        pipe.zremrangebyscore(key, 0, now - window)
        pipe.zadd(key, {str(now): now})
        pipe.zcard(key)
        pipe.expire(key, window + 1)
        results = await pipe.execute()
        count = results[2]
        return count <= max_requests

    async def get_applicable_quotas(self, user_id: str, tenant_id: Optional[str]):
        quotas = []
        quotas.extend(await self._query_quotas("user", user_id))
        if tenant_id:
            quotas.extend(await self._query_quotas("team", tenant_id))
            org_id = await self._get_org_id(tenant_id)
            if org_id:
                quotas.extend(await self._query_quotas("org", org_id))
        return quotas

    async def _query_quotas(self, scope_type, scope_id):
        pass

    async def _get_org_id(self, tenant_id):
        passPython — Rate Limiter

5.4 代理中间件链

# app/proxy/middleware.py

from starlette.middleware.base import BaseHTTPMiddleware

class ProxyMiddleware(BaseHTTPMiddleware):
    async def dispatch(self, request, call_next):
        # Step 1: 认证 — 验证 JWT / API Token
        user = await self._authenticate(request)
        if not user:
            return JSONResponse({"code": 401, "message": "Unauthorized"}, 401)
        request.state.user = user

        # Step 2: 租户解析 — 从 JWT/Header/默认值获取 tenant_id
        tenant_id = await self._resolve_tenant(request, user)
        request.state.tenant_id = tenant_id

        # Step 3: 权限校验 — 检查用户是否有权访问该资源
        if not await self._check_permission(user, request.url.path, request.method):
            return JSONResponse({"code": 403, "message": "Forbidden"}, 403)

        # Step 4: 限流 — 滑动窗口检查
        rate_result = await self._check_rate_limit(user, tenant_id, request.url.path)
        if not rate_result.allowed:
            return JSONResponse(
                {"code": 429, "message": "Rate limit exceeded"},
                429,
                headers={"Retry-After": str(rate_result.retry_after)}
            )

        # Step 5: 审计日志 — 异步记录
        await self._audit_log(user, tenant_id, request)

        # Step 6: 转发请求
        response = await call_next(request)
        return responsePython — Middleware Chain

5.5 主要代理接口清单

RAGFlow 原始接口代理行为限流策略
POST /api/v1/dataset注入 org_id/group_id/space_id 后转发团队级:10次/分钟
POST /api/v1/dataset/{id}/document校验空间权限后转发用户级:30次/分钟
POST /api/v1/dataset/{id}/document/{doc_id}/parse校验解析配额后转发团队级:5次/分钟
POST /api/v1/retrieval注入 tenant_id 过滤后转发用户级:60次/分钟
POST /api/v1/chat/completions校验 Token 配额后转发用户级:30次/分钟
GET /api/v1/dataset按空间权限过滤后转发用户级:120次/分钟
GET /api/v1/user/info直接透传无限流

6 管理服务设计:法人/团队/权限

6.1 管理服务核心职责

管理服务提供 RAGFlow 不具备的扩展管理能力,通过调用 RAGFlow API 完成底层操作:

┌──────────────────────────────────────────────────────────┐ │ 管理服务 (Management Service) │ │ │ │ 法人管理 用户组管理 团队管理 │ │ ┌──────────┐ ┌──────────┐ ┌──────────────┐ │ │ │ 创建法人 │ │ 创建用户组│ │ 创建团队 │ │ │ │ 配置隔离 │ │ 管理成员 │ │ → 调用 RAGFlow│ │ │ │ 级别 │ │ 树形结构 │ │ 创建 Tenant│ │ │ └──────────┘ └──────────┘ │ 邀请成员 │ │ │ │ → 调用 RAGFlow│ │ │ 角色权限 空间管理 │ 邀请 API │ │ │ ┌──────────┐ ┌──────────┐ │ 管理知识库 │ │ │ │ 定义角色 │ │ 创建空间 │ │ → 调用 RAGFlow│ │ │ │ 分配权限 │ │ 授权成员 │ │ KB API │ │ │ │ 用户授权 │ │ 公共/共享 │ └──────────────┘ │ │ └──────────┘ │ /私有 │ │ │ └──────────┘ 限流管理 │ │ ┌──────────────┐ │ │ 文档管理 │ 配置限流规则 │ │ │ ┌──────────┐ │ 查询用量 │ │ │ │ 跨租户 │ │ Token 配额 │ │ │ │ 文档视图 │ └──────────────┘ │ │ │ 批量操作 │ │ │ └──────────┘ │ └──────────────────────────────────────────────────────────┘ │ 调用 RAGFlow API ▼ ┌──────────────────────────────────────────────────────────┐ │ RAGFlow API (不修改) │ │ POST /api/v1/user/register → 创建用户+默认Tenant │ │ POST /api/v1/dataset → 创建知识库 │ │ POST /api/v1/tenant/members → 邀请团队成员 │ │ GET /api/v1/dataset → 列出知识库 │ │ ... │ └──────────────────────────────────────────────────────────┘

6.2 RAGFlow API 客户端

# app/services/ragflow_client.py

import httpx

class RAGFlowClient:
    def __init__(self, base_url: str, api_key: str):
        self.base_url = base_url
        self.api_key = api_key
        self.client = httpx.AsyncClient(
            base_url=base_url,
            headers={"Authorization": f"Bearer {api_key}"},
            timeout=300.0
        )

    async def create_tenant_for_team(self, team_name: str, owner_email: str):
        resp = await self.client.post("/api/v1/user/register", json={
            "email": owner_email,
            "nickname": team_name,
            "password": self._generate_temp_password()
        })
        return resp.json()

    async def invite_member(self, tenant_id: str, email: str):
        resp = await self.client.post("/api/v1/tenant/members", json={
            "email": email
        }, headers={"X-Tenant-ID": tenant_id})
        return resp.json()

    async def create_dataset(self, tenant_id: str, name: str, **kwargs):
        resp = await self.client.post("/api/v1/dataset", json={
            "name": name, **kwargs
        }, headers={"X-Tenant-ID": tenant_id})
        return resp.json()

    async def list_datasets(self, tenant_id: str):
        resp = await self.client.get("/api/v1/dataset",
            headers={"X-Tenant-ID": tenant_id})
        return resp.json()

    async def get_tenant_info(self, tenant_id: str):
        resp = await self.client.get("/api/v1/tenant/info",
            headers={"X-Tenant-ID": tenant_id})
        return resp.json()Python — RAGFlow Client

6.3 团队创建流程(管理服务 → RAGFlow API)

管理API: 创建团队 写入扩展DB 调用 RAGFlow 创建 Tenant 记录 team↔tenant_id 映射 创建默认空间
# app/services/team_service.py

class TeamService:
    def __init__(self, db_session, ragflow_client: RAGFlowClient):
        self.db = db_session
        self.rf = ragflow_client

    async def create_team(self, org_id: str, group_id: str,
                          name: str, owner_email: str) -> Team:
        team = Team(id=uuid4(), org_id=org_id, group_id=group_id, name=name)
        self.db.add(team)

        rf_result = await self.rf.create_tenant_for_team(name, owner_email)
        tenant_id = rf_result["data"]["tenant_id"]

        team.ragflow_tenant_id = tenant_id
        self.db.commit()

        await self._create_default_spaces(team.id, tenant_id, owner_email)

        return teamPython — Team Service

7 数据模型设计

关键设计:扩展服务拥有独立数据库,不修改 RAGFlow 任何表。通过 ragflow_tenant_id 字段建立映射。

7.1 核心表结构

7.1.1 Organization(法人实体表)

CREATE TABLE organization (
    id              VARCHAR(32)  PRIMARY KEY,
    name            VARCHAR(200) NOT NULL,
    code            VARCHAR(64)  UNIQUE NOT NULL,
    description     TEXT,
    logo            TEXT,
    status          CHAR(1)      DEFAULT '1',
    isolation_level VARCHAR(16)  DEFAULT 'logical',
    created_by      VARCHAR(32),
    created_at      DATETIME     DEFAULT CURRENT_TIMESTAMP,
    updated_at      DATETIME     DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
);SQL

7.1.2 UserGroup(用户组表)

CREATE TABLE user_group (
    id              VARCHAR(32)  PRIMARY KEY,
    org_id          VARCHAR(32)  NOT NULL,
    name            VARCHAR(200) NOT NULL,
    code            VARCHAR(64)  NOT NULL,
    parent_id       VARCHAR(32),
    status          CHAR(1)      DEFAULT '1',
    created_by      VARCHAR(32),
    created_at      DATETIME     DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_org_code (org_id, code)
);SQL

7.1.3 Team(团队表)— 核心映射表

CREATE TABLE team (
    id                  VARCHAR(32)  PRIMARY KEY,
    org_id              VARCHAR(32)  NOT NULL,
    group_id            VARCHAR(32)  NOT NULL,
    name                VARCHAR(200) NOT NULL,
    ragflow_tenant_id   VARCHAR(32)  NOT NULL,  -- 映射到 RAGFlow Tenant
    ragflow_owner_email VARCHAR(255),            -- RAGFlow Tenant owner
    status              CHAR(1)      DEFAULT '1',
    created_by          VARCHAR(32),
    created_at          DATETIME     DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_ragflow_tenant (ragflow_tenant_id)
);SQL

7.1.4 OrgUser / UserGroupMember

CREATE TABLE org_user (
    id          VARCHAR(32)  PRIMARY KEY,
    org_id      VARCHAR(32)  NOT NULL,
    user_id     VARCHAR(32)  NOT NULL,
    ragflow_user_id VARCHAR(32),         -- 映射到 RAGFlow User
    role        VARCHAR(32)  DEFAULT 'member',
    status      CHAR(1)      DEFAULT '1',
    created_at  DATETIME     DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_org_user (org_id, user_id)
);

CREATE TABLE user_group_member (
    id          VARCHAR(32)  PRIMARY KEY,
    group_id    VARCHAR(32)  NOT NULL,
    user_id     VARCHAR(32)  NOT NULL,
    role        VARCHAR(32)  DEFAULT 'member',
    status      CHAR(1)      DEFAULT '1',
    created_at  DATETIME     DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_group_user (group_id, user_id)
);SQL

7.1.5 Permission / Role / UserRole

CREATE TABLE permission (
    id            VARCHAR(32)  PRIMARY KEY,
    code          VARCHAR(64)  UNIQUE NOT NULL,
    name          VARCHAR(200) NOT NULL,
    resource_type VARCHAR(32),
    action        VARCHAR(32),
    description   TEXT
);

CREATE TABLE role (
    id          VARCHAR(32)  PRIMARY KEY,
    code        VARCHAR(64)  UNIQUE NOT NULL,
    name        VARCHAR(200) NOT NULL,
    scope       VARCHAR(16) DEFAULT 'team',
    description TEXT
);

CREATE TABLE role_permission (
    role_id       VARCHAR(32) NOT NULL,
    permission_id VARCHAR(32) NOT NULL,
    PRIMARY KEY (role_id, permission_id)
);

CREATE TABLE user_role (
    id          VARCHAR(32)  PRIMARY KEY,
    user_id     VARCHAR(32)  NOT NULL,
    role_id     VARCHAR(32)  NOT NULL,
    scope_type  VARCHAR(16),
    scope_id    VARCHAR(32),
    granted_by  VARCHAR(32),
    created_at  DATETIME DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_user_role_scope (user_id, role_id, scope_type, scope_id)
);SQL

7.1.6 Space(空间表)

CREATE TABLE space (
    id          VARCHAR(32)  PRIMARY KEY,
    team_id     VARCHAR(32)  NOT NULL,
    name        VARCHAR(200) NOT NULL,
    space_type  VARCHAR(16)  NOT NULL,  -- public / shared / private
    description TEXT,
    created_by  VARCHAR(32),
    status      CHAR(1) DEFAULT '1',
    created_at  DATETIME DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_team_space (team_id, name)
);

CREATE TABLE space_member (
    id          VARCHAR(32)  PRIMARY KEY,
    space_id    VARCHAR(32)  NOT NULL,
    user_id     VARCHAR(32)  NOT NULL,
    granted_by  VARCHAR(32),
    created_at  DATETIME DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_space_user (space_id, user_id)
);SQL

7.1.7 RateLimitQuota + AuditLog

CREATE TABLE rate_limit_quota (
    id               VARCHAR(32)  PRIMARY KEY,
    scope_type       VARCHAR(16) NOT NULL,
    scope_id         VARCHAR(32) NOT NULL,
    api_path_pattern VARCHAR(255),
    max_requests     INT NOT NULL,
    window_seconds   INT NOT NULL DEFAULT 60,
    max_tokens       INT,
    max_storage_mb   INT,
    created_at       DATETIME DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY uk_scope_pattern (scope_type, scope_id, api_path_pattern)
);

CREATE TABLE audit_log (
    id          VARCHAR(32)  PRIMARY KEY,
    user_id     VARCHAR(32)  NOT NULL,
    tenant_id   VARCHAR(32),
    action      VARCHAR(64)  NOT NULL,
    resource    VARCHAR(64),
    resource_id VARCHAR(32),
    detail      JSON,
    ip_address  VARCHAR(64),
    created_at  DATETIME DEFAULT CURRENT_TIMESTAMP,
    INDEX idx_user_time (user_id, created_at),
    INDEX idx_tenant_time (tenant_id, created_at)
);SQL

7.2 映射关系图

扩展服务数据库 RAGFlow 数据库(只读引用) ───────────── ───────────────── Organization User │ │ ├── UserGroup UserTenant │ │ │ │ └── Team ── ragflow_tenant_id ────┘→ Tenant │ │ │ │ ├── Space Knowledgebase │ ├── UserRole Document │ └── RateLimitQuota Dialog / Agent │ └── OrgUser ── ragflow_user_id ──────→ User 注意:箭头 → 表示逻辑引用,扩展服务不直接读写 RAGFlow 数据库, 而是通过 RAGFlow API 间接操作。

8 向量数据库租户分区

8.1 分区策略(通过 RAGFlow API 管理)

扩展服务不直接操作向量数据库,而是通过 RAGFlow 的 Dataset API 和配置接口管理分区策略:

策略隔离性实现方式扩展服务职责
索引级分区★★★★★每个法人独立 ES 索引通过 RAGFlow 配置 API 设置索引名规则
Collection 级分区★★★★每个团队一个 Milvus Collection创建团队时调用 RAGFlow Dataset API 初始化
Partition Key 分区★★★共享 Collection + tenant_id代理层注入 tenant_id 元数据过滤
元数据过滤★★RAGFlow Tag Sets代理层注入 org_id/group_id 标签

8.2 代理层向量查询改写

# app/proxy/request_transformer.py

class RequestTransformer:
    async def transform_retrieval_request(self, request_body: dict,
                                           tenant_id: str, org_id: str) -> dict:
        if "metadata_filter" not in request_body:
            request_body["metadata_filter"] = {}
        request_body["metadata_filter"]["tenant_id"] = tenant_id
        request_body["metadata_filter"]["org_id"] = org_id
        return request_body

    async def transform_dataset_create(self, request_body: dict,
                                        team: Team, space: Space) -> dict:
        request_body["meta_fields"] = {
            "org_id": team.org_id,
            "group_id": team.group_id,
            "team_id": team.id,
            "space_type": space.space_type
        }
        return request_bodyPython — Request Transformer

9 空间模型与访问隔离

9.1 空间类型

┌─────────────────────────────────────────────────────┐ │ Team (团队) │ │ ↔ RAGFlow Tenant │ │ │ │ ┌──────────────────────────────────────────────┐ │ │ │ Public Space (公共空间) │ │ │ │ - 所有团队成员自动可见 │ │ │ │ - team_admin 可管理内容 │ │ │ └──────────────────────────────────────────────┘ │ │ │ │ ┌──────────────────────────────────────────────┐ │ │ │ Shared Space (共享空间) │ │ │ │ - 指定成员/角色可见 │ │ │ │ - space_member 表控制 │ │ │ └──────────────────────────────────────────────┘ │ │ │ │ ┌──────────────────────────────────────────────┐ │ │ │ Private Space (私有空间) │ │ │ │ - 仅创建者可见 │ │ │ └──────────────────────────────────────────────┘ │ └─────────────────────────────────────────────────────┘

9.2 空间访问控制(代理层实现)

# app/services/space_service.py

class SpaceService:
    async def filter_accessible_datasets(self, user_id: str,
                                          team_id: str, datasets: list) -> list:
        accessible = []
        for ds in datasets:
            space = await self._get_space_by_dataset(ds["id"], team_id)
            if not space:
                accessible.append(ds)
                continue

            if space.space_type == "public":
                accessible.append(ds)
            elif space.space_type == "shared":
                if await self._is_space_member(user_id, space.id):
                    accessible.append(ds)
            elif space.space_type == "private":
                if ds.get("created_by") == user_id:
                    accessible.append(ds)

        return accessiblePython — Space Access Control

9.3 代理层响应过滤

当用户通过代理层请求 GET /api/v1/dataset 时,代理层在返回 RAGFlow 响应前,根据空间权限过滤:

# app/proxy/response_filter.py

class ResponseFilter:
    async def filter_dataset_list(self, response_body: dict,
                                   user_id: str, team_id: str) -> dict:
        datasets = response_body.get("data", [])
        filtered = await self.space_service.filter_accessible_datasets(
            user_id, team_id, datasets
        )
        response_body["data"] = filtered
        return response_bodyPython — Response Filter

10 API 设计

10.1 扩展服务 API(/api/v1/ext/)

模块方法路径说明
法人管理POST/api/v1/ext/organizations创建法人
GET/api/v1/ext/organizations列出法人
GET/api/v1/ext/organizations/{id}获取法人详情
PUT/api/v1/ext/organizations/{id}更新法人
DELETE/api/v1/ext/organizations/{id}删除法人
用户组POST/api/v1/ext/organizations/{id}/groups创建用户组
GET/api/v1/ext/organizations/{id}/groups列出用户组
POST/api/v1/ext/groups/{id}/members添加组成员
DELETE/api/v1/ext/groups/{id}/members/{uid}移除组成员
团队管理POST/api/v1/ext/teams创建团队(自动创建 RAGFlow Tenant)
GET/api/v1/ext/teams列出我所在的团队
POST/api/v1/ext/teams/{id}/members邀请成员(调用 RAGFlow API)
PUT/api/v1/ext/teams/{id}/members/{uid}/role修改成员角色
DELETE/api/v1/ext/teams/{id}/members/{uid}移除成员
角色权限POST/api/v1/ext/teams/{id}/roles创建自定义角色
PUT/api/v1/ext/teams/{id}/roles/{rid}更新角色权限
POST/api/v1/ext/teams/{id}/members/{uid}/roles分配角色
空间管理POST/api/v1/ext/teams/{id}/spaces创建空间
GET/api/v1/ext/teams/{id}/spaces列出空间
POST/api/v1/ext/teams/{id}/spaces/{sid}/members添加空间成员
限流管理POST/api/v1/ext/rate-limits创建限流规则
GET/api/v1/ext/rate-limits查询限流规则
PUT/api/v1/ext/rate-limits/{id}修改限流规则
GET/api/v1/ext/rate-limits/usage查询当前用量
文档管理GET/api/v1/ext/admin/documents跨租户文档列表
PUT/api/v1/ext/admin/documents/{id}/status修改文档状态
DELETE/api/v1/ext/admin/documents/{id}删除文档

10.2 代理接口(透传到 RAGFlow)

所有非 /api/v1/ext/ 前缀的请求,经代理层鉴权+限流后转发到 RAGFlow:

原始路径代理行为
/api/v1/dataset/*鉴权 → 限流 → 空间权限过滤 → 转发
/api/v1/retrieval鉴权 → 限流 → 注入 tenant_id 过滤 → 转发
/api/v1/chat/*鉴权 → 限流 → Token 配额检查 → 转发
/api/v1/document/*鉴权 → 限流 → 空间权限校验 → 转发
/api/v1/user/*鉴权 → 转发(无限流)
其他鉴权 → 限流 → 转发

11 部署架构

11.1 Docker Compose 部署

# docker/docker-compose.yml

version: "3.8"

services:
  ragflow-extension:
    build: ..
    container_name: ragflow-extension
    ports:
      - "9380:9380"       # 对外端口(替代原 RAGFlow 端口)
    environment:
      - RAGFLOW_INTERNAL_URL=http://ragflow-server:9380
      - RAGFLOW_API_KEY=${RAGFLOW_API_KEY}
      - DATABASE_URL=postgresql+asyncpg://ext:ext@ext-db:5432/ragflow_ext
      - REDIS_URL=redis://redis:6379/1
    depends_on:
      - ext-db
    networks:
      - ragflow-network

  ext-db:
    image: postgres:16-alpine
    container_name: ragflow-ext-db
    environment:
      POSTGRES_DB: ragflow_ext
      POSTGRES_USER: ext
      POSTGRES_PASSWORD: ext
    volumes:
      - ext-db-data:/var/lib/postgresql/data
    networks:
      - ragflow-network

  ragflow-server:
    # RAGFlow 原始服务,端口改为 9381(内部端口)
    # 或保持 9380,扩展服务监听 80/443
    ...

networks:
  ragflow-network:
    external: true   # 复用 RAGFlow 已有网络

volumes:
  ext-db-data:Docker Compose

11.2 端口规划

服务对外端口内部端口说明
ragflow-extension93809380统一入口,客户端只访问此端口
ragflow-server9381仅内部可访问,不对外暴露
ext-db5432扩展服务数据库,仅内部

12 RAGFlow 升级兼容策略

12.1 升级兼容矩阵

RAGFlow 升级场景扩展服务影响应对措施
API 路径不变无影响直接升级 RAGFlow
API 请求/响应字段新增无影响代理层透传,不感知新字段
API 请求/响应字段删除可能影响代理层仅依赖核心字段,非核心字段忽略
API 路径变更需要适配代理层配置路径映射表,无需改代码
认证机制变更需要适配更新代理层认证模块
数据库表结构变更无影响扩展服务不依赖 RAGFlow 数据库
新增 API无影响代理层 catch-all 自动透传

12.2 API 契约测试

# tests/test_ragflow_compat.py

import pytest
from app.services.ragflow_client import RAGFlowClient

RAGFLOW_URL = "http://ragflow-server:9381"

@pytest.mark.asyncio
async def test_dataset_api_compat():
    client = RAGFlowClient(RAGFLOW_URL, "test-key")
    resp = await client.list_datasets(tenant_id="test")
    assert "data" in resp
    assert isinstance(resp["data"], list)

@pytest.mark.asyncio
async def test_retrieval_api_compat():
    client = RAGFlowClient(RAGFLOW_URL, "test-key")
    resp = await client.client.post("/api/v1/retrieval", json={
        "question": "test",
        "dataset_ids": [],
    }, headers={"X-Tenant-ID": "test"})
    assert resp.status_code in (200, 404)

@pytest.mark.asyncio
async def test_chat_api_compat():
    client = RAGFlowClient(RAGFLOW_URL, "test-key")
    resp = await client.client.post("/api/v1/chat/completions", json={
        "model": "test",
        "messages": [{"role": "user", "content": "hi"}],
    }, headers={"X-Tenant-ID": "test"})
    assert resp.status_code in (200, 400, 404)Python — Compatibility Tests

12.3 升级流程

1. 升级 RAGFlow 2. 运行契约测试 3. 适配路径映射 4. 灰度验证 5. 全量切换

13 数据迁移方案

13.1 迁移步骤

Phase 1: 部署扩展服务

  • 部署 ragflow-extension 容器
  • 初始化扩展服务数据库(创建所有扩展表)
  • 配置 RAGFlow API Key 和内部 URL
  • 验证代理层转发正常

Phase 2: 数据同步

  • 通过 RAGFlow API 读取所有现有 Tenant 列表
  • 为每个 Tenant 创建默认 Organization + UserGroup + Team
  • 建立 team.ragflow_tenant_id 映射
  • 同步现有 UserTenant 关系到扩展服务
  • 为现有 owner 角色映射为 team_admin

Phase 3: 流量切换

  • 修改 Nginx/DNS,将外部流量指向扩展服务
  • RAGFlow 原端口改为仅内部访问
  • 灰度开放限流功能
  • 灰度开放空间权限过滤

Phase 4: 验证与优化

  • 全量功能验证
  • 性能测试(代理层延迟测量)
  • 安全审计
  • 全量上线

13.2 数据同步脚本

# scripts/migrate_existing_tenants.py

import asyncio
from app.services.ragflow_client import RAGFlowClient
from app.db.session import get_session
from app.models import Organization, UserGroup, Team, OrgUser

async def migrate():
    rf = RAGFlowClient("http://ragflow-server:9381", "admin-key")
    db = get_session()

    tenants = await rf.list_all_tenants()

    for t in tenants:
        org = Organization(
            id=t["id"], name=f"{t['name']}_org",
            code=f"org_{t['id']}", isolation_level="logical"
        )
        db.add(org)

        group = UserGroup(
            id=t["id"], org_id=org.id,
            name=f"{t['name']}_group", code=f"group_{t['id']}"
        )
        db.add(group)

        team = Team(
            id=t["id"], org_id=org.id, group_id=group.id,
            name=t["name"], ragflow_tenant_id=t["id"]
        )
        db.add(team)

        for member in t.get("members", []):
            if member["role"] == "owner":
                ou = OrgUser(org_id=org.id, user_id=member["user_id"],
                             role="org_admin")
                db.add(ou)

    db.commit()

asyncio.run(migrate())Python — Migration Script

14 实施路线图

Phase 1 — 代理层基础(2 周)

  • FastAPI 项目脚手架
  • 代理路由实现(catch-all 转发)
  • JWT 认证中间件
  • 租户上下文解析
  • Docker Compose 部署

Phase 2 — 限流与权限(2 周)

  • Redis 滑动窗口限流器
  • 限流规则管理 API
  • RBAC 权限引擎
  • 权限管理 API
  • 审计日志

Phase 3 — 扩展管理服务(3 周)

  • 法人/用户组/团队管理 API
  • RAGFlow API 客户端
  • 团队创建 → RAGFlow Tenant 自动化
  • 成员邀请 → RAGFlow 邀请 API 集成
  • 空间模型与访问控制
  • 响应过滤(空间权限)

Phase 4 — 数据迁移与上线(2 周)

  • 现有 Tenant 数据同步
  • 流量切换
  • 契约测试
  • 性能测试
  • 安全审计
  • 全量上线

15 关键风险与应对

风险影响应对措施
代理层增加延迟每次请求多一跳,约 2-5ms连接池复用 + 异步转发 + 热路径缓存
RAGFlow API 变更代理层转发失败契约测试 + 路径映射配置 + 灰度升级
代理层单点故障所有请求不可用多副本部署 + 健康检查 + 降级直连
映射数据不一致team ↔ tenant_id 映射错误定期同步校验 + 对账脚本
RAGFlow 升级大版本API 不兼容契约测试前置 + 版本适配层 + 回滚方案
限流误杀正常用户被限流白名单 + 灰度 + 监控告警

15.1 降级策略

降级方案:当扩展服务不可用时,Nginx 配置 fallback 直接转发到 RAGFlow 原始服务,确保核心功能不受影响:

proxy_pass http://ragflow-extension:9380;
proxy_next_upstream error timeout http_502 http_503;
proxy_next_upstream_tries 2;

15.2 V1 vs V2 方案最终对比

维度V1 源码侵入V2 代理隔离
RAGFlow 源码修改大量修改零修改
升级 RAGFlow 难度高(合并冲突)低(仅验证 API)
额外延迟2-5ms/请求
部署复杂度低(单服务)中(多一个服务)
功能扩展灵活性中(受限于源码)高(独立迭代)
数据隔离可靠性高(直接修改)高(复用原生租户)
长期维护成本
可复用性低(绑定 RAGFlow)高(可适配其他 RAG 引擎)