用 pgai 构建 FastAPI 语义搜索与 RAG 应用:完整实战指南 用 pgai 构建 FastAPI 语义搜索与 RAG 应用完整实战指南【免费下载链接】pgaiA suite of tools to develop RAG, semantic search, and other AI applications more easily with PostgreSQL项目地址: https://gitcode.com/GitHub_Trending/pg/pgai本文基于 pgai 仓库中的 simple_fastapi_app 示例 编写带你完整跑通一个基于 PostgreSQL 的 RAG 应用应用启动时通过pgai.install()安装数据库对象并拉起 Vectorizer 后台 Worker对 Wikipedia 文章列自动创建向量嵌入再用 pgvector 实现语义搜索最后将检索结果作为上下文喂给 LLM 完成 Retrieval Augmented GenerationRAG。读完后你可以掌握 pgai 的核心工作流——建表、建 vectorizer、跟踪嵌入进度、用余弦距离做检索、以及利用触发器实现数据变了、嵌入自动更新并能将其套用到自己的 FastAPI 业务中。1. 示例应用概览该目录包含两个可运行的 FastAPI 应用演示同一套 pgai 能力区别仅在于数据访问方式with_psycopg.py使用psycopg3异步连接池直接写 SQL语义检索基于 pgvector 的操作符with_sqlalchemy.py使用 SQLAlchemy ORM通过pgai.sqlalchemy.vectorizer_relationship把嵌入表自动映射为 ORM 关系。两个应用的功能流程一致启动时安装 pgai 所需的数据库对象并以后台任务形式运行 Vectorizer Worker创建wiki表从 Hugging Face 的wikimedia/wikipedia数据集流式加载少量英文文章为wiki.text列创建 vectorizer自动生成 384 维嵌入Ollamaall-minilm模型提供/vectorizer_status、/search、/insert_pgai_article、/rag四个端点。2. 环境准备与启动2.1 前置条件一个可连接的 PostgreSQL 数据库示例默认连接串为postgresql://postgres:postgreslocalhost:5432/test两个应用文件中都写在该常量处见 with_psycopg.py本地或可达的 Ollama 服务http://localhost:11434并拉取all-minilm嵌入与tinyllama生成两个模型Python 依赖fastapi、psycopg3.x、psycopg_pool、pgvector、ollama、datasets、pgai、numpySQLAlchemy 版本另需sqlalchemy。2.2 获取示例代码并启动在本地克隆仓库后进入示例目录即可直接运行对应 README 中通过 curl 下载with_psycopg.py的做法git clone https://gitcode.com/GitHub_Trending/pg/pgai cd pgai/examples/simple_fastapi_app pip install pgai fastapi uvicorn psycopg psycopg-pool pgvector ollama datasets numpy fastapi dev with_psycopg.py启动后访问http://0.0.0.0:8000/docs即可查看自动生成的 API 文档并在线试调各端点。主要端点为/vectorizer_status查看嵌入创建进度/search?query...语义搜索/insert_pgai_articlePOST向wiki表插入一条关于 pgai 的文章用于演示嵌入自动更新/rag?query...基于检索上下文的 RAG 问答。需要注意的适用前提pgai.install()要求PostgreSQL 15 及以上源码在 install.py 中显式检查server_version_num 15并抛出异常并且会自动执行CREATE EXTENSION IF NOT EXISTS vector即 pgvector。3. 应用启动流程安装 pgai 并拉起 Worker示例应用把数据库初始化 后台任务全部放进 FastAPI 的lifespan生命周期钩子中这是 with_psycopg.py 的核心from contextlib import asynccontextmanager from fastapi import FastAPI from pgai.vectorizer import Worker asynccontextmanager async def lifespan(_app: FastAPI): # 1. 安装 pgai 库到数据库幂等已安装则忽略 pgai.install(DB_URL) # 2. 初始化连接池安装完成后再 open保证 ai schema 对象已就绪 await pool.open() # 3. 以后台 asyncio 任务运行 Vectorizer Worker worker Worker(DB_URL) task asyncio.create_task(worker.run()) # 4. 建表并加载演示数据 await create_wiki_table() if await wiki_table_is_empty(): await load_wiki_articles() # 5. 创建 vectorizer重复启动时忽略 already exists await create_vectorizer() yield # 应用在此运行处理请求 # ---- 关闭阶段优雅停止 Worker ---- print(gracefully shutting down worker...) await worker.request_graceful_shutdown() try: result await asyncio.wait_for(task, timeout20) except asyncio.TimeoutError: print(Worker did not shutdown in time, killing it) app FastAPI(lifespanlifespan)要点解析pgai.install(DB_URL)同步地把内置的 ai.sql 脚本执行到目标库。从 install.py 的实现看它会依次做校验 PG 版本、确保vector扩展存在、读取向量扩展所在 schema、执行安装 SQL若已安装则捕获DuplicateObject错误并静默通过strictFalse默认行为因此可以在应用每次启动时安全调用。生产环境如 README 所建议这一步也可以改用数据库迁移完成而不是放在应用启动路径里。Worker(DB_URL)asyncio.create_task(worker.run())Worker 是嵌入创建的实际执行者。从 worker.py 的构造函数可以看到其可调参数poll_interval默认 1 分钟轮询间隔、once处理完一轮即退出、vectorizer_ids只处理指定 vectorizer缺省为动态模式——自动发现库中所有 vectorizer、concurrency并发数。示例中全部使用默认值即动态发现 每分钟轮询。Worker 部署位置的取舍示例为简化起见把 Worker 跑在 FastAPI 进程内README 明确提示生产上也可以把它放到独立进程或独立容器中运行参见 vectorizer worker 文档。优雅关闭worker.request_graceful_shutdown()会置位内部asyncio.Event见 worker.pyWorker 在当前批次处理完后退出asyncio.wait_for(task, timeout20)给出 20 秒上限兜底。4. 建表与加载 Wikipedia 演示数据4.1 创建wiki表async def create_wiki_table(): async with pool.connection() as conn: async with conn.cursor() as cur: await cur.execute( CREATE TABLE IF NOT EXISTS wiki ( id SERIAL PRIMARY KEY, url TEXT NOT NULL, title TEXT NOT NULL, text TEXT NOT NULL ) ) await conn.commit()表结构很简单自增主键id 三个文本列。其中text将作为 vectorizer 的加载列。4.2 从 Hugging Face 数据集流式加载文章async def load_wiki_articles(): # to keep the demo fast, we have some simple limits num_articles 10 max_text_length 1000 wiki_dataset load_dataset(wikimedia/wikipedia, 20231101.en, splittrain, streamingTrue) async with pool.connection() as conn: async with conn.cursor() as cur: for article in wiki_dataset.take(num_articles): await cur.execute( INSERT INTO wiki (url, title, text) VALUES (%s, %s, %s), (article[url], article[title], article[text][:max_text_length]) ) await conn.commit()这里有两个刻意的演示性限制只取前 10 篇文章、每篇截断到 1000 字符目的是让 demo 快速可跑。streamingTrue让datasets库按迭代器方式逐条拉取数据避免下载整个数据集load_wiki_articles只在wiki_table_is_empty()为真时执行保证重启应用不会重复灌数据。5. 创建 Vectorizer一次声明嵌入自动同步要让wiki.text可被语义检索需要为它生成向量嵌入并保持与数据同步这正是 pgai vectorizer 的职责。示例中的创建代码with_psycopg.pyfrom pgai.vectorizer import CreateVectorizer from pgai.vectorizer.configuration import ( EmbeddingOllamaConfig, LoadingColumnConfig, ) async def create_vectorizer(): vectorizer_statement CreateVectorizer( sourcewiki, target_tablewiki_embedding_storage, loadingLoadingColumnConfig(column_nametext), embeddingEmbeddingOllamaConfig( modelall-minilm, dimensions384, base_urlhttp://localhost:11434 ) ).to_sql() try: async with pool.connection() as conn: async with conn.cursor() as cur: await cur.execute(vectorizer_statement) await conn.commit() except Exception as e: if already exists in str(e): pass # vectorizer 已存在时忽略 else: raise e5.1CreateVectorizer是 SQL 语句构建器CreateVectorizer是 create_vectorizer.py 中定义的 Python 参数模型to_sql()会将其渲染为一条SELECT ai.create_vectorizer(...)调用模板见 config_generator.py。每个配置段loading、embedding、chunking、destination等都继承自 configuration.py 的SQLArgumentMixin各自映射到一个数据库函数例如EmbeddingOllamaConfig对应ai.embedding_ollama(...)。也就是说你在 Python 里写的这段代码最终等价于一条可以直接在 psql 中执行的SELECT ai.create_vectorizer(wiki, loading ai.loading_column(...), embedding ai.embedding_ollama(...), ...)——两条路径Python builder / 原生 SQL创建的是同一个数据库对象完整 SQL API 参考见 vectorizer API reference使用概览见 vectorizer overview。本例涉及的参数及其含义参数取值作用sourcewiki源表regclassvectorizer 监听其主键变化target_tablewiki_embedding_storage嵌入目标表名vectorizer 同时会自动派生检索视图本例为wiki_embedding见第 6 节loading.column_nametext从源表哪个列读取文本embedding.modelall-minilmOllama 嵌入模型名embedding.dimensions384向量维度决定目标表vector(384)列embedding.base_urlhttp://localhost:11434Ollama 服务地址本地部署无需 API Key未显式给出的配置项走数据库侧默认值chunking 默认为none整行作为一个 chunk、formatting、indexing 等也有各自默认。若想调整分块策略比如长文档场景可传入ChunkingCharacterTextSplitterConfig(chunk_size..., chunk_overlap...)等参数全部可选项定义在 configuration.py 与 create_vectorizer.py。5.2 幂等处理应用重启会再次执行create_vectorizer()因此代码捕获了already exists异常并忽略。这是 demo 里的实用手法正式项目建议用 Alembic 迁移管理向量器参见 alembic-integration。6. 跟踪嵌入创建进度ai.vectorizer_status视图Vectorizer 的嵌入创建是异步的Worker 从数据库工作队列中批量取行、调用嵌入服务、再写回目标表。这样设计的目的是支持批量处理并从嵌入服务的瞬时故障中恢复详见 README 第 4 步的说明。要观察进度示例提供了一个只读端点app.get(/vectorizer_status) async def vectorizer_status(): async with pool.connection() as conn: async with conn.cursor(row_factorydict_row) as cur: await cur.execute(SELECT * FROM ai.vectorizer_status) return await cur.fetchall()curl -X GET \ http://0.0.0.0:8000/vectorizer_status \ -H accept: application/jsonai.vectorizer_status是一个由 pgai 安装脚本创建的视图定义在 ai.sql。从视图定义可以看到它返回的列包括vectorizer 的id、name、source_table、target_table、检索view名、embedding_column以及pending_items来自ai.vectorizer_queue_pending(v.id)即工作队列中尚未处理的行数。当pending_items为 0 时说明该 vectorizer 的存量数据已全部完成嵌入——对本 demo 的 10 篇短文这个过程通常只需几秒。从队列的实现还能看到 Worker 的健壮性设计取任务查询使用FOR UPDATE SKIP LOCKED加pg_try_advisory_xact_lock的组合见 vectorizer.py保证多个 Worker 并发安全、互不重复处理同一行批次大小默认 50 行文档类加载为 1可由processing.batch_size覆盖取值被钳制在 1–2048见 vectorizer.py。7. 基于 pgvector 的语义搜索7.1 检索视图与余弦距离操作符Vectoriser 创建后除了目标表wiki_embedding_storage还会生成一个把源表列与嵌入拼在一起的检索视图——本例即wiki_embeddingview_name默认为target_table去掉_storage后缀。该视图包含wiki表的所有列外加embedding向量和chunk该条嵌入对应的文本分片两列。示例中的检索函数with_psycopg.pydataclass class WikiSearchResult: id: int url: str title: str text: str chunk: str distance: float async def _find_relevant_chunks(client: ollama.AsyncClient, query: str, limit: int 2): response await client.embed(modelall-minilm, inputquery) embedding np.array(response.embeddings[0]) async with pool.connection() as conn: async with conn.cursor(row_factoryclass_row(WikiSearchResult)) as cur: await cur.execute( SELECT w.id, w.url, w.title, w.text, w.chunk, w.embedding %s as distance FROM wiki_embedding w ORDER BY distance LIMIT %s , (embedding, limit)) return await cur.fetchall()这段代码的关键点查询向量用与源数据相同的嵌入模型all-minilm生成通过 Ollama SDK 的embed接口获得并转成numpy数组以便 psycopg 以 pgvector 类型绑定连接池在创建时用register_vector_async(conn)注册了向量类型适配器见 with_psycopg.pyembedding %s是 pgvector 定义的余弦距离操作符距离越小越相似所以ORDER BY distance LIMIT n即取最相关的 n 个分片为什么需要 chunking大段文本要拆成语义自洽的小片段每个片段单独嵌入检索时才可能精确命中vectorizer 在创建嵌入时会自动完成切分查询侧直接拿到的是最相关的chunk。示例返回结果中同时带出文章全文text与命中的chunk——不同应用可以按需只取其一README 第 5 步的解释class_row(WikiSearchResult)是 psycopg 的行工厂把每行直接映射进 dataclass/search端点随后用asdict序列化为 JSON。/search端点本身只有三行app.get(/search) async def search(query: str): client ollama.AsyncClient(hosthttp://localhost:11434) results await _find_relevant_chunks(client, query) return [asdict(result) for result in results]试试下面的查询——Properties of Light 这几个词可能根本不出现在任何文章里但嵌入捕捉了语义含义相关段落仍会被排到前面curl -X GET \ http://0.0.0.0:8000/search?queryProperties%20of%20Light \ -H accept: application/json8. 数据变更后嵌入自动更新语义搜索是独立有用的能力也是 RAG 的基石组件。示例用一个端点模拟业务数据变化——向wiki表插入一条关于 pgai 的文章app.post(/insert_pgai_article) async def insert_pgai_article(): async with pool.connection() as conn: async with conn.cursor() as cur: await cur.execute( INSERT INTO wiki (url, title, text) VALUES (%s, %s, %s) , ( https://en.wikipedia.org/wiki/Pgai, pgai - Power your AI applications with PostgreSQL, pgai is a tool to make developing RAG and other AI applications easier... )) await conn.commit() return {message: Article inserted successfully}注意这里没有任何创建嵌入的代码。Vectoriser 在源表上安装了触发器INSERT提交后新行会进入工作队列Worker 在下一轮轮询中自动为其生成嵌入数据更新或删除时对应嵌入也会被同步更新或清理Worker 内部先删旧嵌入再写新嵌入的处理见 vectorizer.py。几秒后用与主题相近的查询再检索就能看到新条目出现在结果中curl -X GET \ http://0.0.0.0:8000/search?queryAI%20Tools \ -H accept: application/json9. 用 RAG 回答 LLM 没见过的问题LLM 没有在 pgai 的资料上训练过靠数据库里的数据才能回答相关问题——这正是 RAG 的价值。/rag端点with_psycopg.py把第 7 节的检索结果组织成提示词上下文再交给 Ollama 上的tinyllama生成回答app.get(/rag) async def rag(query: str) - Optional[str]: # 1. 初始化 Ollama 客户端并检索相关分片 client ollama.AsyncClient(hosthttp://localhost:11434) chunks await _find_relevant_chunks(client, query) # 2. 把检索到的文章拼成上下文 context \n\n.join( f{chunk.title}:\n{chunk.text} for chunk in chunks ) logger.debug(fContext: {context}) # 3. 构造带上下文的提示词 prompt fQuestion: {query} Please use the following context to provide an accurate response: {context} Answer: # 4. 调用 LLM 生成回答 response await client.generate( modeltinyllama, promptprompt, streamFalse ) return response[response]调用效果curl -X GET \ http://0.0.0.0:8000/rag?queryWhat%20is%20pgai \ -H accept: application/json整体链路即用户问题 → 嵌入 → pgvector 余弦检索 top-k 分片 → 拼接为 prompt 上下文 → LLM 生成答案。检索部分完全由 PostgreSQL 承担LLM 只负责基于给定上下文作答这就是 RAG 的分工。10. SQLAlchemy 变体用vectorizer_relationship声明嵌入关系with_sqlalchemy.py 演示了用 ORM 表达同样的流程。模型定义里只需加一行vectorizer_relationshipwith_sqlalchemy.pyclass Wiki(Base): __tablename__ wiki id: Mapped[int] mapped_column(primary_keyTrue) url: Mapped[str] title: Mapped[str] text: Mapped[str] # 为 text 字段声明向量嵌入关系 text_embeddings vectorizer_relationship( target_tablewiki_embeddings, dimensions384 )从实现sqlalchemy/init.py看vectorizer_relationship即_Vectorizer描述符会在 mapper 配置完成后自动动态生成一个嵌入 ORM 模型默认表名源表名_embedding_store或显式指定的target_table含embedding_uuid主键、chunk、embeddingVector(dimensions)列、chunk_seq以及回指父模型的parent关系复制父表主键列并建立ondeleteCASCADE的外键约束把生成的类挂载为Wiki.text_embeddings_model并把 relationship 注册为Wiki.text_embeddings因此可以直接session.query(Wiki).join(Wiki.text_embeddings)。对应的向量检索代码因此可以完全 ORM 化with_sqlalchemy.pyresult session.query( Wiki, Wiki.text_embeddings.embedding.cosine_distance(embedding).label(distance) ).join(Wiki.text_embeddings).order_by( distance ).limit(limit).all()其中cosine_distance(...)是 pgvector 对 SQLAlchemy 的封装pgvector.sqlalchemy.Vector语义与 psycopg 版本里的相同。创建 vectorizer 的调用参数与 psycopg 版完全一致只是经由Session.execute(sqlalchemy.text(...))提交。两个变体怎么选如果应用本来就用 SQLAlchemywith_sqlalchemy.py让嵌入表和源表之间的 join/距离计算都留在 ORM 层代码更统一如果追求轻量或直接用 SQLwith_psycopg.py的路径更短、也更接近 pgai 的数据库对象本身。11. 小结与延伸阅读安装与初始化pgai.install(DB_URL)幂等地安装aischema 下的全部对象要求 PG 15自动装 pgvector实现见 install.pyWorker 生命周期应用内以asyncio.create_task(worker.run())启动、以request_graceful_shutdown()停止生产可拆分为独立进程/容器参考 worker 文档Vectoriser 声明CreateVectorizer(...).to_sql()是SELECT ai.create_vectorizer(...)的 Python 构建器参数模型在 configuration.py 与 create_vectorizer.pySQL 全量参考在 api-reference进度监控ai.vectorizer_status视图的pending_items归零即表示存量嵌入完成视图定义见 ai.sql检索检索视图 pgvector余弦距离 ORDER BY ... LIMIT n分片文本在chunk列数据同步触发器 工作队列FOR UPDATE SKIP LOCKED advisory lock 保证并发安全实现嵌入自动创建/更新/删除队列设计见 vectorizer.pyRAG检索 top-k 分片 → 拼上下文 → LLM 生成端点代码见 with_psycopg.py。想进一步扩展时建议按 vectorizer overview 与 Python 集成文档 了解 chunking、indexing如 HNSW等更多配置项并结合 tests 下的测试用例观察各参数在真实 PostgreSQL 中的行为。【免费下载链接】pgaiA suite of tools to develop RAG, semantic search, and other AI applications more easily with PostgreSQL项目地址: https://gitcode.com/GitHub_Trending/pg/pgai创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考