GitHub周趋势2026W25 | Headroom 压缩 95% Token、NVIDIA 开源 AI Agent 安全扫描器、…
2026-07-28
2026-07-30 0
本文围绕07 | 把字段与指标同步到 Qdrant(生成阶段)整理关键信息和实用建议,帮助读者快速了解主题重点。
07 | 把字段与指标同步到 Qdrant(生成阶段)这是一篇系列文,请按顺序阅读。本文目标NL2SQL 生成 SQL 之前,需要先召回与用户问题相关的字段和指标。上一篇已经把结构化元数据同步进 MySQL:column_info:字段信息
这是一篇系列文,请按顺序阅读。

NL2SQL 生成 SQL 之前,需要先召回与用户问题相关的字段和指标。
上一篇已经把结构化元数据同步进 MySQL:
column_info:字段信息metric_info:指标信息本文继续完成生成阶段的第 2 步:
本文正文只讲这个项目怎样完成同步。向量、向量模型、向量数据库等通用知识,统一放在文末的三个科普模块:
MySQL meta 库├── column_info└── metric_info│▼组装 embedding_textname + description + alias│▼EmbeddingClient → TEI│▼1024 维向量│▼VectorSyncService 组装 Point│▼QdrantRepository → Qdrant├── column_info_collection└── metric_info_collection
各层职责:
| 层 | 文件 | 职责 |
|---|---|---|
| Infrastructure | app/infrastructure/embedding_client.py | 调用 TEI,把文本变成向量 |
| Infrastructure | app/infrastructure/qdrant_client.py | 创建 Qdrant 官方异步客户端 |
| Repository | app/repositories/qdrant_repository.py | 重建 Collection、写入 Point |
| Service | app/services/vector_sync_service.py | 元数据转文本、批量向量化、组装 Point |
| Script | conf/sync_db.py | 读取配置、装配依赖、执行同步、释放资源 |
本文采用“一条元数据记录对应一个向量”,不是给 name、description、alias 分别生成三个向量。
字段记录:
ColumnInfo(id="fact_order.order_amount",name="order_amount",description="订单金额。",alias=["销售额", "订单金额", "收入"],...)
先转换成待向量化输入:
business_id:fact_order.order_amountembedding_text:字段名称:order_amount字段描述:订单金额。字段别名:销售额、订单金额、收入payload:字段 ID、表 ID、类型、角色、别名、样例等原始信息
TEI 生成向量后,最终写入 Qdrant:
{"id": "根据业务 ID 稳定生成的 UUID v5","vector": [0.021, -0.138, 0.407, ...],"payload": {"metadata_type": "column_info","id": "fact_order.order_amount","name": "order_amount","description": "订单金额。","alias": ["销售额", "订单金额", "收入"],"embedding_text": "字段名称:...","embedding_model": "BAAI/bge-large-zh-v1.5",...},}
这里有两个 ID:
id:Qdrant 存储 ID,使用 UUIDpayload["id"]:业务 ID,供应用识别字段或指标conf/app_config.yaml:
qdrant:host: localhostport: 6333embedding_size: 1024column_collection: column_info_collectionmetric_collection: metric_info_collectiontimeout: 60embedding:host: localhostport: 8081model: BAAI/bge-large-zh-v1.5# 当前 CPU TEI 后端最多并行处理 4 条,批次过大会导致请求排队并超时。batch_size: 4timeout: 120
两个服务不要混淆:
| 服务 | 端口 | 负责 |
|---|---|---|
| TEI | 8081 | 把文本转换成向量 |
| Qdrant | 6333 | 保存向量并执行相似度检索 |
依赖:
uv add httpx "qdrant-client==1.16.2"
httpx 用于调用 TEI /embedqdrant-client 是 Qdrant 官方 Python 客户端;这里与 Docker 中的 Qdrant v1.16 保持同一 minor 版本文件:app/infrastructure/embedding_client.py
这个 Client 只负责调用本地 TEI:
class EmbeddingClient:def __init__(self, base_url: str, timeout: float = 60.0):self._client = httpx.AsyncClient(base_url=base_url.rstrip("/"),timeout=timeout,trust_env=False,)async def embed(self, texts: list[str]) -> list[list[float]]:if not texts:return []response = await self._client.post("/embed",json={"inputs": texts,"truncate": True,},)response.raise_for_status()vectors = response.json()if not isinstance(vectors, list) or len(vectors) != len(texts):raise ValueError("Embedding 服务返回数量异常")return vectors
LangChain 常用“文档向量”和“查询向量”两个接口,因此增加两个薄封装:
async def aembed_documents(self,texts: list[str],) -> list[list[float]]:return await self.embed(texts)async def aembed_query(self, text: str) -> list[float]:return (await self.embed([text]))[0]
它们底层都调用同一个 /embed:
aembed_documents ─┐├─ embed() → TEI /embedaembed_query ─────┘
生成索引时调用 aembed_documents();以后用户提问时调用 aembed_query()。
文件底部提供了:
if __name__ == "__main__":asyncio.run(main())
可以直接运行:
uv run python app/infrastructure/embedding_client.py 订单金额 客单价
控制台会显示文本、向量维度和前 8 个值,不打印完整的 1024 维向量。
文件:app/infrastructure/qdrant_client.py
项目使用官方 AsyncQdrantClient。为了让 Client 的构造方式统一,增加一个薄适配类:
from qdrant_client import AsyncQdrantClientclass QdrantClient(AsyncQdrantClient):def __init__(self,base_url: str,timeout: float = 60.0,):super().__init__(url=base_url.rstrip("/"),timeout=int(timeout),)
现在两个 Client 的创建方式具有可预测性:
embedding_client = EmbeddingClient(base_url=embedding_url,timeout=60,)qdrant_client = QdrantClient(base_url=qdrant_url,timeout=60,)
QdrantClient 不重复封装官方方法,仍然可以直接调用:
await qdrant_client.collection_exists(...)await qdrant_client.create_collection(...)await qdrant_client.upsert(...)await qdrant_client.query_points(...)await qdrant_client.close()
文件:app/repositories/qdrant_repository.py
Repository 接收外部创建好的 Client:
class QdrantRepository:def __init__(self, client: AsyncQdrantClient):self._client = client
它不读取配置、不创建 Client,也不关闭 Client。谁创建 Client,谁负责关闭。
async def reset_collection(self,collection_name: str,vector_size: int,) -> None:if await self._client.collection_exists(collection_name):await self._client.delete_collection(collection_name)await self._client.create_collection(collection_name=collection_name,vectors_config=models.VectorParams(size=vector_size,distance=models.Distance.COSINE,),)
当前是生成阶段的全量重建策略:删除旧 Collection,再按 1024 维和 Cosine 距离重建。
async def upsert(self,collection_name: str,points: list[dict[str, Any]],) -> None:qdrant_points = [models.PointStruct(id=point["id"],vector=point["vector"],payload=point.get("payload"),)for point in points]await self._client.upsert(collection_name=collection_name,points=qdrant_points,wait=True,)
upsert 表示:ID 不存在就插入,ID 已存在就更新。
业务 ID 是:
fact_order.order_amountGMV
Qdrant Point ID 使用整数或 UUID,因此通过 UUID v5 稳定转换:
def stable_point_id(collection_name: str,document_id: str,) -> str:return str(uuid5(NAMESPACE_URL,f"{collection_name}:{document_id}",))
相同 Collection 和业务 ID 每次生成相同 UUID,因此同步可重复执行。
文件:app/services/vector_sync_service.py
Service 负责这条业务链路:
ColumnInfo / MetricInfo→ _EmbeddingInput→ 分批调用 TEI→ 检查向量维度→ 组装 Qdrant Point→ 重建 Collection→ 分批 Upsert
字段和指标结构不同,先转换成统一内部类型:
class _EmbeddingInput:business_id: strembedding_text: strpayload: dict[str, Any]
三个字段含义明确:
business_id:用于稳定生成 Point UUIDembedding_text:送给 TEI 生成向量payload:随 Point 写入 Qdrant 的业务信息async def sync_columns(self,collection_name: str,columns: list[ColumnInfo],) -> int:inputs = [ self._column_to_embedding_input(column) for column in columns ]return await self._sync_inputs(collection_name, inputs)async def sync_metrics(self,collection_name: str,metrics: list[MetricInfo],) -> int:inputs = [ self._metric_to_embedding_input(metric) for metric in metrics ]return await self._sync_inputs(collection_name, inputs)
两个入口只负责转换数据,公共同步流程放在 _sync_inputs()。
async def _sync_inputs(self,collection_name: str,inputs: list[_EmbeddingInput],) -> int:points = []# 1. 先生成全部向量,此时不修改旧 Collection。for start in range(0, len(inputs), self.batch_size):batch = inputs[start : start + self.batch_size]vectors = await self.embedding_client.aembed_documents([item.embedding_text for item in batch])self._validate_vectors(vectors)for item, vector in zip(batch, vectors, strict=True):points.append(self._build_point(collection_name, item, vector))# 2. 全部向量成功后,再重建并写入 Collection。await self.qdrant_repository.reset_collection(collection_name,self.vector_size,)for start in range(0, len(points), self.batch_size):await self.qdrant_repository.upsert(collection_name,points[start : start + self.batch_size],)return len(inputs)
为什么先生成全部向量,再删除旧 Collection?
字段文本:
字段名称:order_amount字段描述:订单金额。字段别名:销售额、订单金额、收入
指标文本:
指标名称:GMV指标描述:所有订单的成交金额总和。指标别名:成交总额、订单总额
只有 name + description + alias 参与向量化;类型、角色、所属表和相关字段等结构化信息保留在 payload。
文件:conf/sync_db.py
async def sync_to_qdrant(column_infos: list[ColumnInfo],metric_infos: list[MetricInfo],) -> tuple[int, int]:embedding_config = app_config["embedding"]qdrant_config = app_config["qdrant"]embedding_client = EmbeddingClient(base_url=(f"http://{embedding_config['host']}:"f"{embedding_config['port']}"),timeout=embedding_config.get("timeout", 60),)qdrant_client = QdrantClient(base_url=(f"http://{qdrant_config['host']}:"f"{qdrant_config['port']}"),timeout=qdrant_config.get("timeout", 60),)repository = QdrantRepository(qdrant_client)service = VectorSyncService(embedding_client=embedding_client,qdrant_repository=repository,vector_size=qdrant_config["embedding_size"],model_name=embedding_config["model"],batch_size=embedding_config.get("batch_size", 4),)try:column_count = await service.sync_columns(qdrant_config["column_collection"],column_infos,)metric_count = await service.sync_metrics(qdrant_config["metric_collection"],metric_infos,)return column_count, metric_countfinally:await asyncio.gather(embedding_client.close(),qdrant_client.close(),)
这里是 Composition Root:
确认服务健康:
curl http://localhost:8081/healthcurl http://localhost:6333/healthz
执行:
uv run python conf/sync_db.py
预期看到:
已写入 Qdrant column_info 24 条已写入 Qdrant metric_info 2 条
网页看Qdrant Dashboard:
http://localhost:6333/dashboard#/collections
也可以通过 API 查看少量 Point:
curl http://localhost:6333/collections/column_info_collection/points/scroll -X POST -H 'Content-Type: application/json' -d '{"limit":3,"with_payload":true,"with_vector":false}'
| 现象 | 原因 | 处理 |
|---|---|---|
连接 8081 失败 | TEI 未健康 | 检查 Embedding 容器日志和 /health |
连接 6333 失败 | Qdrant 未启动 | 检查 Docker 和端口映射 |
| 向量维度不是 1024 | 模型与配置不一致 | 核对实际模型和 embedding_size |
| 重复同步出现多条 Point | 使用随机 UUID | 使用业务 ID 派生的 UUID v5 |
| 删除配置后旧 Point 还存在 | 只做 upsert,没有清理旧数据 | 生成阶段重建 Collection |
| 一批请求失败 | 批次或 Token 总量过大 | 调小 embedding.batch_size |
前面的项目结构中,两个 MySQL 数据库分别使用两个模块:
app/dbs/├── dw_db.py└── meta_db.py
它们的代码几乎完全相同,都需要完成:
读取数据库配置→ 创建数据库 URL→ 创建 AsyncEngine→ 创建 Session 工厂→ 对外提供 Session
两个模块真正不同的只有数据库配置:
dw_db:数据仓库,用来读取真实字段类型和样例值meta_db:元数据库,用来保存表、字段、指标及其关系因此没有必要维护两套重复的连接代码。本项目将它们合并成一个通用的 MySQLDatabase:
app/infrastructure/├── embedding_client.py├── mysql_database.py└── qdrant_client.py
infrastructure 是“基础设施”的意思,通常用于存放与外部系统交互的技术实现,例如:
本项目中的职责可以简单理解为:
Service:组织业务流程Repository:表达数据读写操作Infrastructure:建立并管理真实的外部连接
MySQLDatabase 不只是发送一次请求的 Client。它还管理 Engine、连接池、Session 和资源释放,因此使用 Database 比 Client 更准确。
当前只有三个基础设施文件,直接放在 app/infrastructure/ 下更加直观。等以后外部服务明显增多,再拆分 clients/、database/ 等子目录也不迟。
文件:app/infrastructure/mysql_database.py
from sqlalchemy import URLfrom sqlalchemy.ext.asyncio import (AsyncEngine,AsyncSession,async_sessionmaker,create_async_engine,)class MySQLDatabase:def __init__(self,host: str,port: int,user: str,password: str,database: str,*,echo: bool = False,pool_recycle: int = 3600,):database_url = URL.create(drivername="mysql+asyncmy",username=user,password=password,host=host,port=port,database=database,)self.engine: AsyncEngine = create_async_engine(database_url,pool_pre_ping=True,pool_recycle=pool_recycle,echo=echo,connect_args={"charset": "utf8mb4"},)self.session_factory = async_sessionmaker(bind=self.engine,autocommit=False,autoflush=False,expire_on_commit=False,)
这里使用 URL.create(),而不是手工拼接:
f"mysql+asyncmy://{user}:{password}@{host}:{port}/{database}"
这样能够正确处理密码中的 @、:、/ 等 URL 保留字符。
重要参数:
| 参数 | 作用 |
|---|---|
pool_pre_ping=True | 从连接池取连接前先检查连接是否可用 |
pool_recycle=3600 | 定期回收旧连接,避免超过 MySQL wait_timeout |
echo=False | 默认不打印 SQL,调试时可以临时开启 |
charset=utf8mb4 | 支持完整 Unicode,包括 emoji 和生僻字 |
expire_on_commit=False | 提交后保留已加载属性,避免异步访问时意外触发 IO |
三者的关系是:
MySQLDatabase└── AsyncEngine(应用生命周期内复用)└── Connection Pool(管理多条真实连接)├── AsyncSession(一次业务操作)├── AsyncSession(另一次业务操作)└── AsyncSession(另一个请求)
因此不要每查询一个字段就创建一个 Engine,也不要让所有请求长期共用同一个 Session。
MySQLDatabase 提供 session() 上下文管理器:
async def session(self) -> AsyncIterator[AsyncSession]:async with self.session_factory() as session:yield session
同步脚本中的使用方式:
async with meta_database.session() as session, session.begin():...
执行过程:
创建 Session→ 从连接池借用连接→ 开启事务→ 执行业务操作→ 成功时提交,异常时回滚→ 关闭 Session→ 连接归还连接池
类中还提供了 get_session(),以后接入 FastAPI 时可以作为 Depends 使用。session() 和 get_session() 底层共用同一个 Session 工厂。
数据库连接参数继续放在 conf/app_config.yaml 中,MySQLDatabase 本身不读取全局配置。
在 conf/sync_db.py 中创建两个实例:
dw_database = MySQLDatabase(**app_config["dw_db"])meta_database = MySQLDatabase(**app_config["meta_db"])
可以把它理解为:
同一个 MySQLDatabase 类├── dw_db 配置 → dw_database└── meta_db 配置 → meta_database
两个实例分别拥有自己的 Engine 和连接池,但复用了相同的连接管理代码。
同步函数通过参数接收数据库实例:
await sync_to_meta_db(meta_config,dw_database,meta_database,)
这种显式传入的方式有几个好处:
关闭 Session 只是把连接归还连接池,并没有关闭整个连接池。程序结束时还需要释放 Engine:
async def close(self) -> None:await self.engine.dispose()
同步入口使用 finally,保证正常结束或发生异常时都会释放资源:
try:await sync_to_meta_db(meta_config,dw_database,meta_database,)finally:await asyncio.gather(dw_database.close(),meta_database.close(),)
这里遵循统一的生命周期原则:
本模块只讲通用概念,不依赖当前项目代码。
向量可以理解成一组有顺序的数字:
[0.12, -0.37, 0.88]
文本向量是模型对文本特征的数字表示:
“订单金额” → [0.021, -0.138, 0.407, ...]
单个数字通常没有可读业务含义,整组数字共同表达模型学习到的特征。
向量包含多少个数字,就是多少维:
[0.2, 0.5, -0.1] → 3 维
本项目模型输出 1024 个数字,所以是 1024 维。
维度由模型决定,不是越高越好,也不能随意截掉一部分数字。维度越高,存储、传输和计算成本通常也越高。
常见方法:
| 方法 | 留意点 | 常见场景 |
|---|---|---|
| Cosine | 方向是否接近 | 文本语义检索 |
| Dot Product | 点积大小 | 已归一化模型、推荐系统 |
| Euclidean | 空间直线距离 | 数值特征、聚类 |
Cosine 可以先理解成:
方向越接近 → 语义可能越接近
相似度分数不是正确率。0.82 不代表有 82% 的概率正确,阈值需要通过真实问题集评估。
归一化常指把向量长度缩放为 1,同时保留方向:
[2, 0] → [1, 0]
生成文档向量和查询向量时,归一化方式必须一致。更改归一化策略后,通常需要重新生成已有向量。
Embedding 模型把文本、图片或音频转换成向量。本文使用文本模型:
文本 → Tokenizer → 模型 → Pooling → 向量
语义相近的文本通常会得到相近向量,但模型只是提供“相关候选”,并不保证业务判断一定正确。
Tokenizer 会先把文本切成 Token。Token 可能是字、词的一部分、标点或特殊符号。
模型能处理的 Token 数量有上限:
字段和指标描述较短,因此本文使用 truncate: true;长文检索不能只依赖截断,否则尾部信息会丢失。
模型内部通常会为每个 Token 产生表示。为了得到整段文本的单个向量,需要 Pooling:
具体方式由模型训练方式和配置决定,不应随意更改。
语义检索有两个角色:
生成阶段:被搜索内容 → 文档向量 → 存入数据库查询阶段:用户问题 → 查询向量 → 搜索数据库
有些模型要求 Query 和 Document 使用不同前缀:
passage: 订单金额字段query: 销售额是多少
是否需要前缀以及前缀内容,必须查看具体模型说明。不能把其他模型的规则直接照搬过来。
常见考虑因素:
参数量和排行榜只能作为参考,最终要用自己的字段、指标和用户问题评测。
TEI 全称 Text Embeddings Inference,是一个部署文本向量模型的推理服务。
应用 POST /embed↓TEI:分词、批处理、模型推理↓返回向量数组
TEI 负责文本转向量,不负责保存向量或执行相似度检索。
独立部署 TEI 的好处:
常用推理参数与方法:
批次不是越大越好,还受总 Token 数、内存、并发和超时限制。
普通数据库擅长精确条件:
WHERE id = 'GMV'WHERE amount > 100
向量数据库擅长寻找“最相近”的内容:
查询向量 → 找到最相似的 Top K 向量
它不是 MySQL 的完全替代品。常见架构是:
MySQL:保存结构化事实和关系向量数据库:保存向量索引和检索 payload
Collection└── Point├── id├── vector└── payload
| 概念 | 粗略类比 MySQL | 含义 |
|---|---|---|
| Collection | Table | 一类向量的集合 |
| Point | Row | 一条向量记录 |
| Point ID | Primary Key | 唯一标识 |
| Payload | JSON Columns | 可过滤、可返回的业务数据 |
| Vector Index | Index | 加速近邻搜索 |
这只是帮助理解,Qdrant 没有 SQLAlchemy Session,也不是关系型数据库。
数据量大时,逐个比较所有向量成本很高。向量数据库通常使用 Approximate Nearest Neighbor(ANN,近似最近邻)索引。
HNSW 是常见 ANN 算法,通过多层图结构快速找到相近向量。它牺牲少量绝对精确性,换取更快查询速度。
常见调优目标:
向量相似度可以和结构化过滤组合:
语义相似+ metadata_type = column_info+ table_id = fact_order
经常用于过滤的 payload 字段可以建立 Payload Index,避免查看大量 Point。
Upsert:
ID 不存在 → 插入ID 已存在 → 更新
同步任务应使用稳定 ID。随机 UUID v4 每次都不同,重复执行可能产生重复 Point;由业务 ID 派生的 UUID v5 每次相同,更适合幂等同步。
QdrantClient 会复用底层 HTTP 或 gRPC 连接,可以被多个 Repository 和请求共享。
它没有 MySQL ORM Session 的直接对应物:
MySQL:Engine → 多个 Session → 事务Qdrant:一个 Client → 多个并发 API 请求
不要每查询一条数据就创建一个新 Client;应用关闭或一次性任务结束时再统一关闭。
Qdrant 没有 SQLAlchemy 那种跨多个 API 操作的事务回滚。
开发阶段可以:
删除旧 Collection → 重建 → 写入
生产环境更稳妥的方式是:
创建临时 Collection→ 写入全部 Point→ 验证数量和查询→ 切换 Alias→ 删除旧 Collection
这样 Embedding 或写入失败时,线上仍然使用旧索引。
limit:返回多少条候选score_threshold:最低相似度阈值with_payload:是否返回业务数据with_vector:是否返回原始向量阈值和 Top K 需要通过真实业务问题集调优,不能只看单个演示结果。
本文完成了字段和指标向量索引的生成。下一篇继续把维度值同步到 Elasticsearch,之后再实现 Agent 中的:
recall_columnrecall_metricrecall_value以上内容可作为基础参考,实际处理时再结合具体场景灵活调整。