From f724a7a120b32d86fe950f8153e83793b4545a05 Mon Sep 17 00:00:00 2001 From: leonyangdev <2443992009@qq.com> Date: Sat, 29 Aug 2026 12:13:44 +0800 Subject: [PATCH 1/6] feat(v2): implement document intelligence pipeline --- .env.example | 2 + README.md | 58 ++-- apps/api/app.py | 50 ++- apps/api/routes.py | 8 +- apps/api/schemas.py | 40 ++- apps/web/app/chat/page.tsx | 6 +- apps/web/app/knowledge-bases/[id]/page.tsx | 17 +- apps/web/app/lib.ts | 24 ++ apps/web/components/retrieval-evidence.tsx | 7 +- apps/web/package-lock.json | 4 +- apps/web/package.json | 2 +- docker-compose.yml | 3 + docs/4.v2_implementation.md | 296 +++++++++++++++++ pyproject.toml | 9 +- scripts/smoke_v2.py | 276 ++++++++++++++++ src/ultimate_rag/__init__.py | 4 +- src/ultimate_rag/application/context.py | 5 +- src/ultimate_rag/application/services.py | 52 +-- src/ultimate_rag/chunkers/__init__.py | 4 +- src/ultimate_rag/chunkers/markdown.py | 42 +-- src/ultimate_rag/config.py | 9 +- src/ultimate_rag/domain/models.py | 83 ++++- src/ultimate_rag/domain/ports.py | 10 +- .../infrastructure/database/repository.py | 2 +- .../infrastructure/storage/minio.py | 2 +- src/ultimate_rag/ocr/__init__.py | 5 + src/ultimate_rag/ocr/bailian.py | 91 ++++++ src/ultimate_rag/parsers/__init__.py | 15 +- src/ultimate_rag/parsers/_shared.py | 100 ++++++ src/ultimate_rag/parsers/html.py | 139 ++++++++ src/ultimate_rag/parsers/image.py | 94 ++++++ src/ultimate_rag/parsers/markdown.py | 10 +- src/ultimate_rag/parsers/office.py | 253 +++++++++++++++ src/ultimate_rag/parsers/pdf.py | 144 +++++++++ src/ultimate_rag/parsers/registry.py | 2 +- src/ultimate_rag/vectorstores/milvus.py | 19 +- tests/unit/test_bailian_ocr.py | 66 ++++ tests/unit/test_chunker.py | 25 +- tests/unit/test_ingestion_service.py | 6 +- tests/unit/test_milvus_vector_store.py | 22 ++ tests/unit/test_v2_parsers.py | 232 +++++++++++++ uv.lock | 306 +++++++++++++++++- 42 files changed, 2422 insertions(+), 122 deletions(-) create mode 100644 docs/4.v2_implementation.md create mode 100644 scripts/smoke_v2.py create mode 100644 src/ultimate_rag/ocr/__init__.py create mode 100644 src/ultimate_rag/ocr/bailian.py create mode 100644 src/ultimate_rag/parsers/_shared.py create mode 100644 src/ultimate_rag/parsers/html.py create mode 100644 src/ultimate_rag/parsers/image.py create mode 100644 src/ultimate_rag/parsers/office.py create mode 100644 src/ultimate_rag/parsers/pdf.py create mode 100644 tests/unit/test_bailian_ocr.py create mode 100644 tests/unit/test_v2_parsers.py diff --git a/.env.example b/.env.example index b1bacb0..82ca1cc 100644 --- a/.env.example +++ b/.env.example @@ -4,6 +4,8 @@ DASHSCOPE_API_KEY=replace-me EMBEDDING_MODEL=text-embedding-v4 EMBEDDING_DIMENSION=1024 LLM_MODEL=qwen-plus +OCR_MODEL=qwen-vl-ocr-latest +OCR_MAX_IMAGE_BYTES=6291456 # Local Docker Compose defaults DATABASE_URL=postgresql+asyncpg://ultimate_rag:ultimate_rag@localhost:5432/ultimate_rag diff --git a/README.md b/README.md index dc2d76d..ebbde25 100644 --- a/README.md +++ b/README.md @@ -1,21 +1,23 @@ # UltimateRAG -一个从最小可用 RAG 持续演进为企业级知识平台的学习型工程。当前仓库实现 **V1.0 · Naive RAG**: -它既能作为 RAG 全链路学习项目,也保留了真实企业系统需要的数据边界、失败状态、可替换端口和可测试性。 +一个从最小可用 RAG 持续演进为企业级知识平台的学习型工程。当前仓库实现 +**V2.0 · Document Intelligence**:在保留 V1 可运行 RAG 闭环的基础上,把多种原始格式统一为 +可追溯的文档领域模型,使新增 Parser 不需要修改 RAG 主流程。 -## V1 能做什么 +## V2 能做什么 用户可以在 Web 中完成以下闭环: 1. 创建知识库 -2. 上传 UTF-8 Markdown -3. 查看文档从 `PENDING` 到 `READY` 的处理结果 -4. 使用 Milvus Dense Retrieval 独立调试召回内容和分数 -5. 使用阿里云百炼模型进行知识库问答 -6. 查看答案引用的文档、章节和 Chunk -7. 删除文档或知识库,并同步清理三类存储 +2. 上传 Markdown、PDF、DOCX、XLSX、PPTX、HTML 或常见图片 +3. 自动识别 PDF 原生文本页与扫描页,并使用阿里云百炼 Qwen-OCR 处理扫描内容 +4. 查看文档从 `PENDING` 到 `READY` 的处理结果和实际 Parser +5. 使用 Milvus Dense Retrieval 独立调试召回内容和分数 +6. 使用阿里云百炼模型进行知识库问答 +7. 查看答案引用的章节、PDF 页码、Excel 区域或 PPT 幻灯片 +8. 删除文档或知识库,并同步清理三类存储 -V1 明确不包含 PDF、OCR、混合检索、Reranker、Agent、ACL、异步任务和 RAGOps;这些属于后续版本。 +V2 明确不包含混合检索、Reranker、Agent、ACL、异步任务和 RAGOps;这些属于后续版本。 ## 架构 @@ -32,8 +34,9 @@ Application Services │ ▼ Domain Ports - ├── DocumentParser → MarkdownParser - ├── Chunker → StructureAwareMarkdownChunker + ├── DocumentParser → Markdown / PDF / Office / HTML / Image OCR + ├── OCRClient → BailianOCRClient + ├── Chunker → StructureAwareChunker ├── Embedder → BailianEmbedder ├── VectorStore → MilvusVectorStore ├── ObjectStorage → MinioObjectStorage @@ -47,12 +50,13 @@ Domain Ports - Domain 不依赖 FastAPI、SQLAlchemy、Milvus、OpenAI SDK 或 LangChain - PostgreSQL 保存知识库、文档状态和 Chunk 元数据 -- MinIO 保存原始 Markdown,且对象键由系统生成 +- MinIO 保存所有原始文件,且对象键由系统生成 - Milvus 只保存可重建向量索引,不作为业务事实数据源 - 文档仅在 Parse、Chunk、Embedding、Index 全部成功后进入 `READY` - 知识库内容按不可信输入处理,不能覆盖系统 Prompt -详细设计见 [V1 实现说明](docs/3.v1_implementation.md)。 +详细设计见 [V2 实现说明](docs/4.v2_implementation.md),V1 的基础闭环见 +[V1 实现说明](docs/3.v1_implementation.md)。 ## 技术栈 @@ -62,6 +66,7 @@ Domain Ports - 阿里云百炼 OpenAI 兼容 API - Embedding 默认 `text-embedding-v4`,1024 维 - LLM 默认 `qwen-plus` + - OCR 默认 `qwen-vl-ocr-latest` - Next.js 16、React 19、TypeScript、Tailwind CSS 4、shadcn/ui、AI SDK - uv、pytest、Ruff、Mypy @@ -79,6 +84,8 @@ DASHSCOPE_API_KEY=你的API-Key EMBEDDING_MODEL=text-embedding-v4 EMBEDDING_DIMENSION=1024 LLM_MODEL=qwen-plus +OCR_MODEL=qwen-vl-ocr-latest +OCR_MAX_IMAGE_BYTES=6291456 # 可选;留空时浏览器自动访问当前页面主机的 8000 端口 NEXT_PUBLIC_API_URL= @@ -201,7 +208,7 @@ PENDING → PARSING → CHUNKING → EMBEDDING → INDEXING → READY └→ FAILED(任一处理阶段失败) ``` -失败文档保留原文件与错误状态,方便定位问题和未来重建。V1 是同步管线,因此上传请求会等待处理完成。 +失败文档保留原文件与错误状态,方便定位问题和未来重建。V2 仍是同步管线,因此上传请求会等待处理完成。 ## 验证 @@ -218,7 +225,8 @@ npm run build npm audit ``` -固定 Smoke Test 文档位于 `tests/fixtures/rag.md`,推荐问题是“BGE-M3 是什么?”。 +单元测试会在内存中生成各类格式,避免提交二进制 Fixture。固定 Markdown Smoke Test 文档位于 +`tests/fixtures/rag.md`,推荐问题是“BGE-M3 是什么?”。 启动 Docker 全栈后,执行真实 PostgreSQL、MinIO、Milvus 和百炼闭环验收: @@ -226,8 +234,14 @@ npm audit uv run python scripts/smoke_v1.py --api-url http://localhost:8000 ``` -脚本会创建临时知识库、上传 Fixture、验证文档 `READY`、检索命中、流式答案和 Citation, -最后删除临时知识库及其跨存储资源。 +V1 脚本保留用于回归。V2 全格式验收使用: + +```bash +uv run python scripts/smoke_v2.py --api-url http://localhost:8000 +``` + +V2 脚本会动态生成并上传全部支持格式,验证 Parser、`READY`、带来源位置的检索、流式答案和 +Citation,最后删除临时知识库及其跨存储资源。 ## 目录 @@ -236,7 +250,8 @@ apps/web/ Next.js Web apps/api/ FastAPI 应用 src/ultimate_rag/domain/ 领域模型与端口 src/ultimate_rag/application/ 显式业务工作流 -src/ultimate_rag/parsers/ Markdown 解析与注册表 +src/ultimate_rag/parsers/ Markdown / PDF / Office / HTML / Image 解析与注册表 +src/ultimate_rag/ocr/ 百炼 OCR 适配器 src/ultimate_rag/chunkers/ 结构感知切块 src/ultimate_rag/embeddings/ 百炼向量适配器 src/ultimate_rag/vectorstores/ Milvus 适配器 @@ -250,9 +265,10 @@ docs/ 产品、架构与实现文档 ## 安全提醒 -- 上传文件必须是 UTF-8 Markdown,最大 10 MB +- 上传文件最大 10 MB;Markdown/HTML 必须使用 UTF-8,Office 会检查 ZIP Bomb 风险 +- 图片提交 OCR 前会验证真实编码;PDF 最多 500 页,扫描页按页调用 OCR - 用户文件名不参与本地路径或对象键构造 - `.env`、API Key 和生产凭据禁止提交 - 默认 Docker 密码只适合本地开发 - 检索内容和 LLM 输出都视为不可信数据 -- 对公网部署前仍需要认证、ACL、限流与审计;这些不属于 V1 范围 +- 对公网部署前仍需要认证、ACL、限流与审计;这些不属于 V2 范围 diff --git a/apps/api/app.py b/apps/api/app.py index 155de9e..3960e0c 100644 --- a/apps/api/app.py +++ b/apps/api/app.py @@ -9,7 +9,7 @@ 业务顺序由 Application Service 编排,外部协议细节由 Infrastructure Adapter 负责。 设计背景: - V1 使用一个显式 Container 保存进程内共享依赖。相比在每个 Route 中临时创建客户端, + V2 使用一个显式 Container 保存进程内共享依赖。相比在每个 Route 中临时创建客户端, 这种方式可以复用连接池并让对象生命周期可见;当前依赖数量有限,因此不引入 DI 框架。 典型使用场景: @@ -17,7 +17,7 @@ 注意事项 / 已知限制: 数据库 Schema 只能通过 Alembic Migration 管理,启动过程不会自动建表。MinIO Bucket - 或 Milvus Collection 初始化失败时应用不会进入请求服务阶段。V1 关闭时显式释放数据库 + 或 Milvus Collection 初始化失败时应用不会进入请求服务阶段。V2 关闭时显式释放数据库 Engine;其他 Adapter 当前没有统一的异步关闭端口。 """ @@ -38,7 +38,7 @@ RAGService, RetrievalService, ) -from ultimate_rag.chunkers import StructureAwareMarkdownChunker +from ultimate_rag.chunkers import StructureAwareChunker from ultimate_rag.config import get_settings from ultimate_rag.domain.exceptions import ( InvalidDocumentError, @@ -49,7 +49,17 @@ from ultimate_rag.generation import BailianLLMClient from ultimate_rag.infrastructure.database import create_database from ultimate_rag.infrastructure.storage import MinioObjectStorage -from ultimate_rag.parsers import MarkdownParser, ParserRegistry +from ultimate_rag.ocr import BailianOCRClient +from ultimate_rag.parsers import ( + ExcelParser, + HtmlParser, + ImageOCRParser, + MarkdownParser, + ParserRegistry, + PDFParser, + PowerPointParser, + WordParser, +) from ultimate_rag.vectorstores import MilvusVectorStore settings = get_settings() @@ -61,7 +71,7 @@ @asynccontextmanager async def lifespan(app: FastAPI) -> AsyncIterator[None]: - """在 FastAPI 进程生命周期内创建、校验并暴露 V1 所需依赖。 + """在 FastAPI 进程生命周期内创建、校验并暴露 V2 所需依赖。 Lifespan 的启动部分严格先于 ``yield`` 执行。只有 MinIO Bucket 和 Milvus Collection 均准备完成后,FastAPI 才开始接收请求;因此 Route 不需要处理“依赖尚未初始化”的状态。 @@ -110,20 +120,40 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: dimension=settings.embedding_dimension, ) - # LLM 与 Embedder 共享百炼的 OpenAI-Compatible Endpoint 和 API Key,但模型职责不同: - # Embedder 只生成检索向量,LLMClient 只根据构造后的知识上下文生成最终答案。 + # LLM、Embedder 与 OCR 共享百炼 Endpoint 和 API Key,但模型职责完全分离。 llm = BailianLLMClient( api_key=settings.dashscope_api_key, base_url=settings.dashscope_base_url, model=settings.llm_model, timeout=settings.model_timeout_seconds, ) + ocr = BailianOCRClient( + api_key=settings.dashscope_api_key, + base_url=settings.dashscope_base_url, + model=settings.ocr_model, + max_image_bytes=settings.ocr_max_image_bytes, + timeout=settings.model_timeout_seconds, + ) # 阶段 3:装配不直接拥有外部资源的领域策略和应用服务。 # Registry 隔离源格式与 Parser 选择,Chunker 负责统一 ParsedDocument 之后的切块; # RetrievalService 复用同一个 Embedder,保证查询向量与文档向量处于相同向量空间。 - registry = ParserRegistry([MarkdownParser()]) - chunker = StructureAwareMarkdownChunker(settings.chunk_max_chars, settings.chunk_overlap_chars) + registry = ParserRegistry( + [ + MarkdownParser(), + WordParser(), + ExcelParser(), + PowerPointParser(), + HtmlParser(), + PDFParser( + ocr, + native_text_threshold=settings.pdf_native_text_threshold, + render_scale=settings.pdf_render_scale, + ), + ImageOCRParser(ocr), + ] + ) + chunker = StructureAwareChunker(settings.chunk_max_chars, settings.chunk_overlap_chars) retrieval = RetrievalService(embedder, vector_store) # 阶段 4:把已经装配好的对象集中放入进程级 Container。 @@ -162,7 +192,7 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: await engine.dispose() -app = FastAPI(title=settings.app_name, version="1.0.0", lifespan=lifespan) +app = FastAPI(title=settings.app_name, version="2.0.0", lifespan=lifespan) app.add_middleware( CORSMiddleware, allow_origins=settings.cors_origins, diff --git a/apps/api/routes.py b/apps/api/routes.py index 473d081..1e09868 100644 --- a/apps/api/routes.py +++ b/apps/api/routes.py @@ -1,4 +1,4 @@ -"""UltimateRAG V1 HTTP 路由。 +"""UltimateRAG V2 HTTP 路由。 模块职责: 验证 HTTP 输入、调用应用服务,并把领域结果映射为普通 JSON 或 AI SDK UI Message Stream。 @@ -133,13 +133,13 @@ async def upload_document( request: Request, file: Annotated[UploadFile, File()], ) -> DocumentResponse: - """上传并同步完成 Markdown 的解析、切块、向量化和索引。""" + """上传并同步完成多格式文档的解析、切块、向量化和索引。""" dependencies = container(request) content = await _read_bounded_upload(file, dependencies.max_upload_bytes) value = await dependencies.ingestion.ingest( knowledge_base_id, - file.filename or "document.md", - file.content_type or "text/markdown", + file.filename or "document.bin", + file.content_type or "application/octet-stream", content, ) return DocumentResponse.from_domain(value) diff --git a/apps/api/schemas.py b/apps/api/schemas.py index 833bd43..2a940aa 100644 --- a/apps/api/schemas.py +++ b/apps/api/schemas.py @@ -7,7 +7,39 @@ from pydantic import BaseModel, ConfigDict, Field -from ultimate_rag.domain.models import Citation, Document, KnowledgeBase, RetrievalResult +from ultimate_rag.domain.models import ( + Citation, + Document, + KnowledgeBase, + RetrievalResult, + SourceLocator, +) + + +class SourceLocatorResponse(BaseModel): + """跨格式原文位置;不同文档类型只填写适用字段。""" + + heading_path: list[str] = Field(default_factory=list) + page: int | None = None + bbox: list[float] | None = None + sheet: str | None = None + cell_range: str | None = None + slide: int | None = None + + @classmethod + def from_domain(cls, value: SourceLocator | None) -> "SourceLocatorResponse | None": + """把不可变 Locator 映射为 JSON 友好的 API 结构。""" + + if value is None: + return None + return cls( + heading_path=list(value.heading_path), + page=value.page, + bbox=list(value.bbox) if value.bbox else None, + sheet=value.sheet, + cell_range=value.cell_range, + slide=value.slide, + ) class KnowledgeBaseCreate(BaseModel): @@ -93,6 +125,7 @@ class RetrievalResultResponse(BaseModel): filename: str content: str heading_path: list[str] + locator: SourceLocatorResponse | None score: float @classmethod @@ -104,6 +137,7 @@ def from_domain(cls, value: RetrievalResult) -> "RetrievalResultResponse": filename=value.filename, content=value.content, heading_path=list(value.heading_path), + locator=SourceLocatorResponse.from_domain(value.locator), score=value.score, ) @@ -122,6 +156,7 @@ class CitationResponse(BaseModel): filename: str chunk_id: str heading_path: list[str] + locator: SourceLocatorResponse | None @classmethod def from_domain(cls, value: Citation) -> "CitationResponse": @@ -131,11 +166,12 @@ def from_domain(cls, value: Citation) -> "CitationResponse": filename=value.filename, chunk_id=value.chunk_id, heading_path=list(value.heading_path), + locator=SourceLocatorResponse.from_domain(value.locator), ) class ChatResponse(BaseModel): - """答案、引用和基础检索调试信息的完整 V1 响应。""" + """答案、跨格式引用和基础检索调试信息的完整 V2 响应。""" answer: str citations: list[CitationResponse] diff --git a/apps/web/app/chat/page.tsx b/apps/web/app/chat/page.tsx index 6adfd8c..57906ae 100644 --- a/apps/web/app/chat/page.tsx +++ b/apps/web/app/chat/page.tsx @@ -263,7 +263,7 @@ export default function ChatPage() {

{pageError ? "请检查后端服务是否可用,刷新页面重试。" - : "先创建一个知识库并上传 Markdown 文档,然后回到这里开始问答。"} + : "先创建一个知识库并上传文档,然后回到这里开始问答。"}