Skip to content

03 元数据知识库构建

学习理念:本章开发一个构建元数据知识库的脚本(build_meta_knowledge.py),核心是理解三层架构(脚本入口 → 业务服务层 → Repository 持久层)和三存储写入(MySQL 结构化 + Qdrant 向量索引 + ES 全文索引)。

海外对标:Airflow DAG 的数据管道思路、Databricks Unity Catalog 元数据采集

本节 AI 替代率:~65% | 人工干预率:~35%

角色能力范围
🤖 AI 擅长生成实体类、ORM Model、Mapper 转换代码、Repository CRUD
👤 人类需理解三层架构的分层动机、Entity↔Model↔Mapper 的转换链、向量索引构建的"嵌入—存储—检索"全流程

📌 原文说明:以下内容来源于原始笔记第5章"元数据知识库",包含需求说明、代码组织规划、实体类/Model/Model转换、配置、入口脚本、核心业务逻辑、Repository 层、测试。原文全部保留,补充了代码阅读优先级标注和架构说明。

一、需求说明

本章旨在开发一个用于构建元数据知识库的脚本。该脚本以一个 YAML 配置文件作为输入参数,配置文件中定义了需要同步至元数据知识库的表格信息与指标信息。脚本功能包括:

  • 根据配置内容,将指定的表格及指标信息同步至 Meta 数据库;
  • 为字段信息和指标信息构建向量索引;
  • 为字段取值建立全文索引。

image-20260207142744497


二、代码组织规划

构建元数据知识库模块的核心代码采用分层设计,具体组织结构如下:

🔥 【P0 必须理解】 三层架构梳理:

  • 脚本层scripts/build_meta_knowledge.py):命令行入口,argparse 接收配置文件路径,初始化所有客户端,组装依赖
  • 业务层services/meta_knowledge_service.py):核心业务逻辑,6 个私有方法串联成完整构建流程
  • 持久层repositories/):封装 MySQL/Qdrant/ES 的读写,SQLAlchemy ORM 操作 + Qdrant PointStruct + ES bulk

实体层entities/)和 Model 层models/)之间通过 Mapper 类 做类型转换:Entity → to_model() → MySQL Model → to_entity() → Entity

bash
data-agent/
├── conf/
   └── meta_config.yaml # 配置文件,用于指定待同步的表格和指标
└── app/
    ├── conf/
   └── meta_config.py # 用于定义meta_config.yaml的参数结构
    ├── scripts/
   └── build_meta_knowledge.py # 构建元数据知识库的入口脚本,负责解析命令行参数
    ├── services/
   └── meta_knowledge_service.py # 负责实现构建元数据知识库的核心逻辑
    ├── entities/
   ├── table_info.py # 用于定义表格信息的业务实体类
   ├── column_info.py # 用于定义字段信息的业务实体类
   ├── metric_info.py # 用于定义指标信息的业务实体类
   ├── column_metric.py # 用于定义字段指标关系的业务实体类
   └── value_info.py # 用于定义字段取值的业务实体类
    └── repositories/
        ├── mysql/
   ├── meta/
   ├── meta_mysql_repository.py # 负责实现meta数据库的读写操作
   └── mappers/
       ├── table_info_mapper.py
       ├── column_info_mapper.py
       ├── metric_info_mapper.py
       └── column_metric_mapper.py
   └── dw/
       └── dw_mysql_repository.py # 负责实现dw数据库的读写操作
        ├── qdrant/
   ├── column_qdrant_repository.py # 负责实现qdrant中column相关集合的读写操作
   └── metric_qdrant_repository.py # 负责实现qdrant中metric相关集合的读写操作
        └── es/
            └── value_es_repository.py # 用于定义es中value索引的实体类

三、实体层、Model 层与 Mapper

3.1 业务实体类(Entities)

🟡 【P1 看注释就行】 实体类用 @dataclass 定义,职责是"统一封装不同存储系统的数据格式",作为业务层和持久层之间的标准交换格式。

业务实体类用于统一封装不同存储系统中的数据结构(如 Meta 数据库、Elasticsearch、Qdrant),作为应用上层组件进行读写操作时的标准参数类型和返回值类型。本项目定义的业务实体有:TableInfo、ColumnInfo、MetricInfo、ColumnMetric、ValueInfo,具体代码如下:

  • TableInfo(data-agent/app/entities/table_info.py

    python
    from dataclasses import dataclass
    
    @dataclass
    class TableInfo:
        id: str
        name: str
        role: str
        description: str
  • ColumnInfo(data-agent/app/entities/column_info.py

    python
    from dataclasses import dataclass
    from typing import Any
    
    @dataclass
    class ColumnInfo:
        id: str
        name: str
        type: str
        role: str
        examples: list[Any]
        description: str
        alias: list[str]
        table_id: str
  • MetricInfo(data-agent/app/entities/metric_info.py

    python
    from dataclasses import dataclass
    
    @dataclass
    class MetricInfo:
        id: str
        name: str
        description: str
        relevant_columns: list[str]
        alias: list[str]
  • ColumnMetric(data-agent/app/entities/column_metric.py

    python
    from dataclasses import dataclass
    
    @dataclass
    class ColumnMetric:
        column_id: str
        metric_id: str
  • ValueInfo(data-agent/app/entities/value_info.py

    python
    from dataclasses import dataclass
    
    @dataclass
    class ValueInfo:
        id: str
        value: str
        column_id: str

3.2 ORM Model 类

🟡 【P1 看注释就行】 SQLAlchemy 2.0 的 Mapped + mapped_column 声明式语法。每个 Model 类映射一张 MySQL 表。

Python 类 ↔ 数据库表 的映射关系(类 = 表,实例 = 行,属性 = 列),用面向对象的方式操作数据库,无需编写原生 SQL(底层由 SQLAlchemy 自动转换)官方术语Declarative Model(声明式模型),日常简化为 Model

  • data-agent/app/models/base.py

    python
    from sqlalchemy.orm import DeclarativeBase
    
    class Base(DeclarativeBase):
        pass
  • data-agent/app/models/table_info_mysql.py

    python
    from sqlalchemy import String, Text
    from sqlalchemy.orm import Mapped, mapped_column
    from app.models.base import Base
    
    class TableInfoMySQL(Base):
        __tablename__ = "table_info"
        id: Mapped[str] = mapped_column(String(64), primary_key=True, comment="表编号")
        name: Mapped[str | None] = mapped_column(String(128), comment="表名称")
        role: Mapped[str | None] = mapped_column(String(32), comment="表类型(fact/dim)")
        description: Mapped[str | None] = mapped_column(Text, comment="表描述")
  • data-agent/app/models/column_info_mysql.py

    python
    from sqlalchemy import String, Text
    from sqlalchemy.types import JSON
    from sqlalchemy.orm import Mapped, mapped_column
    from app.models.base import Base
    
    class ColumnInfoMySQL(Base):
        __tablename__ = "column_info"
        id: Mapped[str] = mapped_column(String(64), primary_key=True, comment="列编号")
        name: Mapped[str | None] = mapped_column(String(128), comment="列名称")
        type: Mapped[str | None] = mapped_column(String(64), comment="数据类型")
        role: Mapped[str | None] = mapped_column(String(32), comment="列类型(primary_key,foreign_key,measure,dimension)")
        examples: Mapped[dict | list | None] = mapped_column(JSON, comment="数据示例")
        description: Mapped[str | None] = mapped_column(Text, comment="列描述")
        alias: Mapped[dict | list | None] = mapped_column(JSON, comment="列别名")
        table_id: Mapped[str | None] = mapped_column(String(64), comment="所属表编号")
  • data-agent/app/models/metric_info_mysql.py

    python
    from sqlalchemy import String, Text
    from sqlalchemy.types import JSON
    from sqlalchemy.orm import Mapped, mapped_column
    from app.models.base import Base
    
    class MetricInfoMySQL(Base):
        __tablename__ = "metric_info"
        id: Mapped[str] = mapped_column(String(64), primary_key=True, comment="指标编码")
        name: Mapped[str | None] = mapped_column(String(128), comment="指标名称")
        description: Mapped[str | None] = mapped_column(Text, comment="指标描述")
        relevant_columns: Mapped[dict | list | None] = mapped_column(JSON, comment="关联字段")
        alias: Mapped[dict | list | None] = mapped_column(JSON, comment="指标别名")
  • data-agent/app/models/column_metric_mysql.py

    python
    from sqlalchemy import String
    from sqlalchemy.orm import Mapped, mapped_column
    from app.models.base import Base
    
    class ColumnMetricMySQL(Base):
        __tablename__ = "column_metric"
        column_id: Mapped[str] = mapped_column(String(64), primary_key=True, comment="列编号")
        metric_id: Mapped[str] = mapped_column(String(64), primary_key=True, comment="指标编号")

3.3 Mapper 转换类

🟡 【P1 看注释就行】 Mapper 负责 Entity ↔ Model 的双向转换。to_entity() 手动赋值(因为类型不完全一致),to_model()asdict() 快速转换。

  • TableInfoMapper (data-agent/app/repositories/mysql/meta/mappers/table_info_mapper.py

    python
    from dataclasses import asdict
    from app.entities.table_info import TableInfo
    from app.models.table_info_mysql import TableInfoMySQL
    
    class TableInfoMapper:
        @staticmethod
        def to_entity(table_info_mysql: TableInfoMySQL) -> TableInfo:
            return TableInfo(
                id=table_info_mysql.id,
                name=table_info_mysql.name,
                role=table_info_mysql.role,
                description=table_info_mysql.description
            )
    
        @staticmethod
        def to_model(table_info: TableInfo) -> TableInfoMySQL:
            return TableInfoMySQL(**asdict(table_info))
  • ColumnInfoMapper(data-agent/app/repositories/mysql/meta/mappers/column_info_mapper.py

    python
    from dataclasses import asdict
    from app.entities.column_info import ColumnInfo
    from app.models.column_info_mysql import ColumnInfoMySQL
    
    class ColumnInfoMapper:
        @staticmethod
        def to_entity(column_info_mysql: ColumnInfoMySQL) -> ColumnInfo:
            # 将持久层对象转为业务层对象,只能手动赋值
            return ColumnInfo(
                id=column_info_mysql.id,
                name=column_info_mysql.name,
                type=column_info_mysql.type,
                role=column_info_mysql.role,
                examples=column_info_mysql.examples,
                description=column_info_mysql.description,
                alias=column_info_mysql.alias,
                table_id=column_info_mysql.table_id,
            )
        @staticmethod
        def to_model(column_info: ColumnInfo) -> ColumnInfoMySQL:
            # 将业务层对象转为持久层对象
            return ColumnInfoMySQL(**asdict(column_info))
  • MetricInfoMapper(data-agent/app/repositories/mysql/meta/mappers/metric_info_mapper.py

    python
    from dataclasses import asdict
    from app.entities.metric_info import MetricInfo
    from app.models.metric_info_mysql import MetricInfoMySQL
    
    class MetricInfoMapper:
        @staticmethod
        def to_entity(model: MetricInfoMySQL) -> MetricInfo:
            return MetricInfo(
                id=model.id, name=model.name, description=model.description,
                relevant_columns=model.relevant_columns, alias=model.alias
            )
        @staticmethod
        def to_model(entity: MetricInfo):
            return MetricInfoMySQL(**asdict(entity))
  • ColumnMetricMapper(data-agent/app/repositories/mysql/meta/mappers/column_metric_mapper.py

    python
    from dataclasses import asdict
    from app.entities.column_metric import ColumnMetric
    from app.models.column_metric_mysql import ColumnMetricMySQL
    
    class ColumnMetricMapper:
        @staticmethod
        def to_entity(column_metric_mysql: ColumnMetricMySQL):
            return ColumnMetric(column_id=column_metric_mysql.column_id, metric_id=column_metric_mysql.metric_id)
        @staticmethod
        def to_model(column_metric: ColumnMetric):
            return ColumnMetricMySQL(**asdict(column_metric))

四、配置文件

4.1 配置结构定义

🟡 【P1 看注释就行】 MetaConfig 定义了 YAML 配置的 dataclass schema,通过 OmegaConf 加载和校验。

data-agent/app/conf/meta_config.py中编写如下内容:

python
from dataclasses import dataclass
from typing import Optional

@dataclass
class ColumnConfig:
    name: str
    role: str
    description: str
    alias: list[str]
    sync: bool

@dataclass
class TableConfig:
    name: str
    role: str
    description: str
    columns: list[ColumnConfig]

@dataclass
class MetricConfig:
    name: str
    description: str
    relevant_columns: list[str]
    alias: list[str]

@dataclass
class MetaConfig:
    tables: Optional[list[TableConfig]] = None
    metrics: Optional[list[MetricConfig]] = None

4.2 具体配置内容(完整版)

当前数据仓库模拟环境(dw库)中共有 5 张表和多个指标需要同步到元数据知识库。具体配置如下:

**说明:**该配置文件应由脚本调用方提供,所以该文件在项目中的路径不做要求,现暂将其置于data-agent/conf/meta_config.yaml中。

🔥 【P0 必须理解】 sync: true/false 控制该字段的取值是否需要同步到 ES 建立全文索引。只有维度表的业务可读字段(如 province、category、member_level)才需要,主键/外键/度量不需要。

yaml
# data-agent/conf/meta_config.yaml
tables:
  - name: dim_region
    role: dim
    description: 地区维度表,用于描述订单发生的地理区域信息。
    columns:
      - name: region_id
        role: primary_key
        description: 地区唯一标识。
        alias: [ 地区ID, 区域ID ]
        sync: false
      - name: province
        role: dimension
        description: 订单所属的省份名称。
        alias: [ 省份, , 所在省份 ]
        sync: true
      - name: region_name
        role: dimension
        description: 订单所属的大区名称,如华东、华南等。
        alias: [ 地区, 区域, 大区 ]
        sync: true
      - name: country
        role: dimension
        description: 地区所属国家名称。
        alias: [ 国家, 国家名称 ]
        sync: true

  - name: dim_customer
    role: dim
    description: 客户维度表,描述下单客户的基本属性。
    columns:
      - name: customer_id
        role: primary_key
        description: 客户唯一标识。
        alias: [ 客户ID, 用户ID ]
        sync: false
      - name: customer_name
        role: dimension
        description: 客户名称。
        alias: [ 客户名称, 用户名称 ]
        sync: true
      - name: gender
        role: dimension
        description: 客户性别。
        alias: [ 性别 ]
        sync: true
      - name: member_level
        role: dimension
        description: 客户会员等级。
        alias: [ 会员等级, 用户等级 ]
        sync: true

  - name: dim_product
    role: dim
    description: 商品维度表,描述商品的基本属性信息。
    columns:
      - name: product_id
        role: primary_key
        description: 商品唯一标识。
        alias: [ 商品ID, 产品ID ]
        sync: false
      - name: product_name
        role: dimension
        description: 商品名称。
        alias: [ 商品名称, 产品名称 ]
        sync: true
      - name: category
        role: dimension
        description: 商品所属品类。
        alias: [ 商品类别, 品类, 分类 ]
        sync: true
      - name: brand
        role: dimension
        description: 商品品牌名称。
        alias: [ 品牌, 品牌名称 ]
        sync: true

  - name: dim_date
    role: dim
    description: 时间维度表,用于多时间粒度分析。
    columns:
      - name: date_id
        role: primary_key
        description: 日期唯一标识,格式 yyyyMMdd。
        alias: [ 日期ID, 日期 ]
        sync: false
      - name: year
        role: dimension
        description: 年份。
        alias: [ , 年份 ]
        sync: false
      - name: quarter
        role: dimension
        description: 季度。
        alias: [ 季度 ]
        sync: true
      - name: month
        role: dimension
        description: 月份。
        alias: [ , 月份 ]
        sync: false
      - name: day
        role: dimension
        description: 日。
        alias: [ ,  ]
        sync: false

  - name: fact_order
    role: fact
    description: 订单事实表,记录订单数量和金额等核心指标。
    columns:
      - name: order_id
        role: primary_key
        description: 订单唯一标识。
        alias: [ 订单ID ]
        sync: false
      - name: customer_id
        role: foreign_key
        description: 关联客户维度的外键。
        alias: [ 客户ID, 用户ID ]
        sync: false
      - name: product_id
        role: foreign_key
        description: 关联商品维度的外键。
        alias: [ 商品ID, 产品ID ]
        sync: false
      - name: date_id
        role: foreign_key
        description: 关联时间维度的外键。
        alias: [ 日期, 下单日期 ]
        sync: false
      - name: region_id
        role: foreign_key
        description: 关联地区维度的外键。
        alias: [ 地区ID, 区域ID ]
        sync: false
      - name: order_quantity
        role: measure
        description: 订单中商品的购买数量。
        alias: [ 销量, 购买数量, 件数 ]
        sync: false
      - name: order_amount
        role: measure
        description: 订单金额。
        alias: [ 销售额, 订单金额, 收入 ]
        sync: false

metrics:
  - name: GMV
    description: 全称Gross Merchandise Value,表示所有订单的成交金额总和。
    relevant_columns:
      - fact_order.order_amount
    alias: [ 成交总额, 订单总额 ]
  - name: AOV
    description: 全称Average Order Value,表示所有订单的成交金额平均值。
    relevant_columns:
      - fact_order.order_quantity
    alias: [ 平均单价, 平均订单金额 ]
  - name: ORDER_COUNT
    description: 统计周期内的订单总数量,按order_id去重计数。
    relevant_columns:
      - fact_order.order_id
    alias: [ 订单数, 下单量 ]
  - name: ORDER_COUNT_BY_REGION
    description: 按地区维度(region_id/region_name)统计的订单数量。
    relevant_columns:
      - fact_order.order_id
      - dim_region.region_id
      - dim_region.region_name
    alias: [ 区域订单量, 各地区下单数 ]
  - name: SALES_QUANTITY
    description: 统计周期内所有订单的商品购买数量总和。
    relevant_columns:
      - fact_order.order_quantity
    alias: [ 总销量, 销售件数 ]
  - name: SALES_QUANTITY_BY_CATEGORY
    description: 按商品品类(category)统计的商品销售数量。
    relevant_columns:
      - fact_order.order_quantity
      - dim_product.category
    alias: [ 品类销售件数, 分类销量 ]
  - name: AVG_ORDER_QUANTITY
    description: 订单商品购买数量的平均值(总销量/订单总数)。
    relevant_columns:
      - fact_order.order_quantity
      - fact_order.order_id
    alias: [ 单均商品数, 平均购买件数 ]
  - name: AVG_AMOUNT_BY_CUSTOMER
    description: 统计周期内客户的平均订单金额(GMV/客户数)。
    relevant_columns:
      - fact_order.order_amount
      - fact_order.customer_id
    alias: [ 客均消费, 客户平均下单金额 ]
  - name: GMV_BY_MEMBER_LEVEL
    description: 按客户会员等级(member_level)统计的成交总额。
    relevant_columns:
      - fact_order.order_amount
      - dim_customer.member_level
    alias: [ 会员等级销售额, 各等级成交总额 ]
  - name: DAILY_GMV
    description: 按日期维度(date_id/year/month/day)统计的每日GMV。
    relevant_columns:
      - fact_order.order_amount
      - dim_date.date_id
      - dim_date.year
      - dim_date.month
      - dim_date.day
    alias: [ 日GMV, 每日成交总额 ]
  - name: QUARTERLY_SALES
    description: 按季度(quarter)统计的商品销售总量。
    relevant_columns:
      - fact_order.order_quantity
      - dim_date.quarter
    alias: [ 季度销售件数, 季度销量 ]
  - name: GMV_BY_BRAND
    description: 按商品品牌(brand)统计的成交总额。
    relevant_columns:
      - fact_order.order_amount
      - dim_product.brand
    alias: [ 品牌GMV, 各品牌销售额 ]
  - name: TOP_PRODUCT_SALES
    description: 按商品名称(product_name)排序的商品销售数量TOP榜。
    relevant_columns:
      - fact_order.order_quantity
      - dim_product.product_id
      - dim_product.product_name
    alias: [ 商品销量TOP, 热销商品件数 ]
  - name: CUSTOMER_COUNT
    description: 统计周期内下单的客户总数量,按customer_id去重计数。
    relevant_columns:
      - fact_order.customer_id
    alias: [ 下单客户数, 客户数量 ]
  - name: CUSTOMER_COUNT_BY_GENDER
    description: 按客户性别(gender)统计的下单客户数量。
    relevant_columns:
      - fact_order.customer_id
      - dim_customer.gender
    alias: [ 性别客户数, 男女下单人数 ]
  - name: MEMBER_LEVEL_COUNT
    description: 按会员等级(member_level)统计的下单客户数量。
    relevant_columns:
      - fact_order.customer_id
      - dim_customer.member_level
    alias: [ 会员等级客户数, 各等级下单人数 ]

五、入口脚本

🔥 【P0 必须理解】 入口脚本使用 argparse 接收 -c/--conf 参数指定配置路径,初始化所有 5 个客户端,实例化 6 个 Repository,然后调用 MetaKnowledgeService.build() 执行完整构建流程。

data-agent/app/scripts/build_meta_knowledge.py中编写如下代码:

python
import asyncio
from argparse import ArgumentParser
from pathlib import Path

from app.clients.embedding_client_manager import embedding_client_manager
from app.clients.es_client_manager import es_client_manager
from app.clients.mysql_client_manager import (
    meta_mysql_client_manager,
    dw_mysql_client_manager,
)
from app.clients.qdrant_client_manager import qdrant_client_manager
from app.repositories.qdrant.column_qdrant_repository import ColumnQdrantRepository
from app.repositories.mysql.dw.dw_mysql_repository import DWMySQLRepository
from app.repositories.mysql.meta.meta_mysql_repository import MetaMySQLRepository
from app.repositories.es.value_es_repository import ValueESRepository
from app.repositories.qdrant.metric_qdrant_repository import MetricQdrantRepository
from app.services.meta_knowledge_service import MetaKnowledgeService


async def build(config_path: Path):
    meta_mysql_client_manager.init()  # 初始化元数据MySQL客户端
    dw_mysql_client_manager.init()  # 初始化数据仓库MySQL客户端
    qdrant_client_manager.init()  # 初始化Qdrant客户端
    embedding_client_manager.init()  # 初始化Embedding客户端
    es_client_manager.init()  # 初始化Elasticsearch客户端

    async with (
        meta_mysql_client_manager.session_factory() as meta_session,
        dw_mysql_client_manager.session_factory() as dw_session,
    ):
        meta_mysql_repository = MetaMySQLRepository(meta_session)
        dw_mysql_repository = DWMySQLRepository(dw_session)
        column_qdrant_repository = ColumnQdrantRepository(qdrant_client_manager.client)
        embedding_client = embedding_client_manager.client
        value_es_repository = ValueESRepository(es_client_manager.client)
        metric_qdrant_repository = MetricQdrantRepository(qdrant_client_manager.client)

        meta_knowledge_service = MetaKnowledgeService(
            meta_mysql_repository=meta_mysql_repository,
            dw_mysql_repository=dw_mysql_repository,
            column_qdrant_repository=column_qdrant_repository,
            embedding_client=embedding_client,
            value_es_repository=value_es_repository,
            metric_qdrant_repository=metric_qdrant_repository,
        )
        await meta_knowledge_service.build(config_path)

    await meta_mysql_client_manager.close()
    await dw_mysql_client_manager.close()
    await qdrant_client_manager.close()
    await es_client_manager.close()


if __name__ == "__main__":
    parser = ArgumentParser()
    parser.add_argument("-c", "--conf")
    args = parser.parse_args()
    config_path = Path(args.conf)
    asyncio.run(build(config_path))

六、核心业务逻辑(完整版)

核心业务逻辑位于data-agent/app/services/meta_knowledge_service.py中,具体内容如下:

python
import uuid
from pathlib import Path

from langchain_huggingface import HuggingFaceEndpointEmbeddings
from omegaconf import OmegaConf

from app.conf.meta_config import MetaConfig
from app.core.log import logger
from app.entities.column_info import ColumnInfo
from app.entities.column_metric import ColumnMetric
from app.entities.metric_info import MetricInfo
from app.entities.table_info import TableInfo
from app.entities.value_info import ValueInfo
from app.repositories.es.value_es_repository import ValueESRepository
from app.repositories.mysql.dw.dw_mysql_repository import DWMySQLRepository
from app.repositories.mysql.meta.meta_mysql_repository import MetaMySQLRepository
from app.repositories.qdrant.column_qdrant_repository import ColumnQdrantRepository
from app.repositories.qdrant.metric_qdrant_repository import MetricQdrantRepository


class MetaKnowledgeService:
    def __init__(
        self,
        meta_mysql_repository: MetaMySQLRepository,
        dw_mysql_repository: DWMySQLRepository,
        column_qdrant_repository: ColumnQdrantRepository,
        embedding_client: HuggingFaceEndpointEmbeddings,
        value_es_repository: ValueESRepository,
        metric_qdrant_repository: MetricQdrantRepository,
    ):
        self.meta_mysql_repository = meta_mysql_repository
        self.dw_mysql_repository = dw_mysql_repository
        self.column_qdrant_repository = column_qdrant_repository
        self.embedding_client = embedding_client
        self.value_es_repository = value_es_repository
        self.metric_qdrant_repository = metric_qdrant_repository

    async def _save_tables_to_meta_db(self, meta_config: MetaConfig) -> list[ColumnInfo]:
        table_infos: list[TableInfo] = []
        column_infos: list[ColumnInfo] = []

        for table in meta_config.tables:
            # 构造TableInfo实例
            table_info = TableInfo(
                id=table.name,
                name=table.name,
                role=table.role,
                description=table.description,
            )
            table_infos.append(table_info)

            # 查询该表的所有字段类型
            column_types: dict[str, str] = await self.dw_mysql_repository.get_column_types(table.name)
            for column in table.columns:
                # 查询该字段的部分取值作为示例
                column_values: list = await self.dw_mysql_repository.get_column_values(table.name, column.name, 10)
                # 构造ColumnInfo实例
                column_info = ColumnInfo(
                    id=f"{table.name}.{column.name}",
                    name=column.name,
                    type=column_types[column.name],
                    role=column.role,
                    examples=column_values,
                    description=column.description,
                    alias=column.alias,
                    table_id=table.name,
                )
                column_infos.append(column_info)

        # 保存表信息和字段信息到元数据数据库
        async with self.meta_mysql_repository.session.begin():
            await self.meta_mysql_repository.save_table_infos(table_infos)
            await self.meta_mysql_repository.save_column_infos(column_infos)

        return column_infos

    async def _save_column_info_to_qdrant(self, column_infos: list[ColumnInfo]):
        # 确保column_info的collection存在
        await self.column_qdrant_repository.ensure_collection()
        # 构造待保存的数据
        points: list[dict] = []
        for column_info in column_infos:
            points.append(
                {
                    "id": uuid.uuid4(),
                    "embedding_text": column_info.name,
                    "payload": column_info,
                }
            )
            points.append(
                {
                    "id": uuid.uuid4(),
                    "embedding_text": column_info.description,
                    "payload": column_info,
                }
            )
            for alia in column_info.alias:
                points.append(
                    {"id": uuid.uuid4(), "embedding_text": alia, "payload": column_info}
                )
        # 向量列表
        embedding_texts = [point["embedding_text"] for point in points]
        embedding_batch_size = 10
        embeddings = []
        for i in range(0, len(embedding_texts), embedding_batch_size):
            batch_embedding_texts = embedding_texts[i: i + embedding_batch_size]
            batch_embeddings = await self.embedding_client.aembed_documents(
                batch_embedding_texts
            )
            embeddings.extend(batch_embeddings)

        # id列表
        ids = [point["id"] for point in points]
        # payload列表
        payloads = [point["payload"] for point in points]
        # 保存数据到qdrant
        await self.column_qdrant_repository.upsert(ids, embeddings, payloads)

    async def _save_value_info_to_es(self, meta_config: MetaConfig, column_infos: list[ColumnInfo]):
        # 确保index存在
        await self.value_es_repository.ensure_index()

        # 获取需要同步取值的列
        column2sync: dict[str, bool] = {}
        for table in meta_config.tables:
            for column in table.columns:
                column2sync[f"{table.name}.{column.name}"] = column.sync

        # 构造ValueInfo列表
        value_infos: list[ValueInfo] = []
        for column_info in column_infos:
            sync = column2sync[column_info.id]
            if sync:
                # 查询这个列的所有取值
                table_name = column_info.table_id
                column_name = column_info.name
                values = await self.dw_mysql_repository.get_column_values(
                    table_name, column_name, 100000
                )
                current_value_infos = [
                    ValueInfo(
                        id=f"{column_info.id}.{value}",
                        value=value,
                        column_id=column_info.id,
                    )
                    for value in values
                ]
                value_infos.extend(current_value_infos)
        # 批量保存到Elasticsearch
        await self.value_es_repository.index(value_infos)

    async def _save_metrics_to_meta_db(self, meta_config):
        metric_infos: list[MetricInfo] = []
        column_metrics: list[ColumnMetric] = []
        for metric in meta_config.metrics:
            # 构造MetricInfo数据
            metric_info = MetricInfo(
                id=metric.name,
                name=metric.name,
                description=metric.description,
                relevant_columns=metric.relevant_columns,
                alias=metric.alias,
            )
            metric_infos.append(metric_info)
            for relevant_column in metric.relevant_columns:
                # 构造ColumnMetric数据
                column_metric = ColumnMetric(
                    column_id=relevant_column, metric_id=metric.name
                )
                column_metrics.append(column_metric)
        # 保存到元数据数据库
        async with self.meta_mysql_repository.session.begin():
            await self.meta_mysql_repository.save_metric_infos(metric_infos)
            await self.meta_mysql_repository.save_column_metrics(column_metrics)
        return metric_infos

    async def _save_metric_info_to_qdrant(self, metric_infos: list[MetricInfo]):
        # 确保collection存在
        await self.metric_qdrant_repository.ensure_collection()
        # 构造待保存的数据
        points: list[dict] = []
        for metric_info in metric_infos:
            points.append(
                {"id": uuid.uuid4(), "embedding_text": metric_info.name, "payload": metric_info}
            )
            points.append(
                {"id": uuid.uuid4(), "embedding_text": metric_info.description, "payload": metric_info}
            )
            for alia in metric_info.alias:
                points.append(
                    {"id": uuid.uuid4(), "embedding_text": alia, "payload": metric_info}
                )
        ids = [point["id"] for point in points]
        embeddings = []
        embedding_texts = [point["embedding_text"] for point in points]
        embedding_batch_size = 10
        for i in range(0, len(embedding_texts), embedding_batch_size):
            batch_embedding_texts = embedding_texts[i: i + embedding_batch_size]
            batch_embeddings = await self.embedding_client.aembed_documents(batch_embedding_texts)
            embeddings.extend(batch_embeddings)
        payloads = [point["payload"] for point in points]
        # 保存数据到qdrant
        await self.metric_qdrant_repository.upsert(ids, embeddings, payloads)

    async def build(self, config_path: Path):
        # 1.加载配置文件
        context = OmegaConf.load(config_path)
        schema = OmegaConf.structured(MetaConfig)
        meta_config: MetaConfig = OmegaConf.to_object(OmegaConf.merge(schema, context))
        logger.info("加载配置文件")
        # 2.处理表信息
        if meta_config.tables:
            # 2.1 保存表信息到meta数据库
            column_infos = await self._save_tables_to_meta_db(meta_config)
            logger.info("保存表信息到meta数据库")
            # 2.2 为字段信息建立向量索引
            await self._save_column_info_to_qdrant(column_infos)
            logger.info("为字段信息建立向量索引")
            # 2.3 为字段取值建立全文索引
            await self._save_value_info_to_es(meta_config, column_infos)
            logger.info("为字段取值建立全文索引")
        # 3.处理指标信息
        if meta_config.metrics:
            # 3.1 保存指标信息到meta数据库
            metric_infos = await self._save_metrics_to_meta_db(meta_config)
            logger.info("保存指标信息到meta数据库")
            # 3.2 为指标信息建立向量索引
            await self._save_metric_info_to_qdrant(metric_infos)
            logger.info("为指标信息建立向量索引")
        logger.info("元数据知识库构建完成")

七、Repository 持久层

7.1 meta_mysql_repository

meta_mysql_repository用于实现meta数据库的读写操作,具体代码如下:

python
from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession
from app.entities.column_info import ColumnInfo
from app.entities.column_metric import ColumnMetric
from app.entities.metric_info import MetricInfo
from app.entities.table_info import TableInfo
from app.models.column_info_mysql import ColumnInfoMySQL
from app.models.table_info_mysql import TableInfoMySQL
from app.repositories.mysql.meta.mappers.column_info_mapper import ColumnInfoMapper
from app.repositories.mysql.meta.mappers.column_metric_mapper import ColumnMetricMapper
from app.repositories.mysql.meta.mappers.metric_info_mapper import MetricInfoMapper
from app.repositories.mysql.meta.mappers.table_info_mapper import TableInfoMapper


class MetaMySQLRepository:
    def __init__(self, session: AsyncSession):
        self.session = session

    async def save_table_infos(self, table_infos: list[TableInfo]):
        models = [TableInfoMapper.to_model(table_info) for table_info in table_infos]
        self.session.add_all(models)

    async def save_column_infos(self, columns_info: list[ColumnInfo]):
        models = [ColumnInfoMapper.to_model(column_info) for column_info in columns_info]
        self.session.add_all(models)

    async def save_metric_infos(self, metric_infos: list[MetricInfo]):
        self.session.add_all([MetricInfoMapper.to_model(metric_info) for metric_info in metric_infos])

    async def save_column_metrics(self, column_metrics: list[ColumnMetric]):
        self.session.add_all([ColumnMetricMapper.to_model(column_metric) for column_metric in column_metrics])

    async def get_column_info_by_id(self, column_id: str) -> ColumnInfo | None:
        result: ColumnInfoMySQL | None = await self.session.get(ColumnInfoMySQL, column_id)
        if result:
            return ColumnInfoMapper.to_entity(result)
        return None

    async def get_table_info_by_id(self, table_id: str) -> TableInfo | None:
        result: TableInfoMySQL | None = await self.session.get(TableInfoMySQL, table_id)
        if result:
            return TableInfoMapper.to_entity(result)
        return None

    async def get_key_columns_by_table_id(self, table_id: str) -> list[ColumnInfo]:
        sql = """
            select * from column_info
            where table_id = :table_id
            and role in ('primary_key', 'foreign_key')
        """
        result = await self.session.execute(text(sql), {"table_id": table_id})
        return [ColumnInfo(**row) for row in result.mappings().fetchall()]

7.2 dw_mysql_repository

dw_mysql_repository用于实现模拟数据仓库(dw库)的读写操作。

python
from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession


class DWMySQLRepository:
    def __init__(self, session: AsyncSession):
        self.session = session

    async def get_column_types(self, table_name: str) -> dict[str, str]:
        sql = f"show columns from {table_name}"
        result = await self.session.execute(text(sql))
        return {row.Field: row.Type for row in result.fetchall()}

    async def get_column_values(self, table_name: str, column_name: str, limit: int):
        sql = f"select distinct {column_name} from {table_name} limit {limit}"
        result = await self.session.execute(text(sql))
        return result.scalars().fetchall()

    async def get_db_info(self):
        result = await self.session.execute(text("select version()"))
        version = result.scalar()
        dialect = self.session.get_bind().dialect.name
        return {'version': version, 'dialect': dialect}

    async def validate_sql(self, sql):
        await self.session.execute(text(f"explain {sql}"))

    async def execute_sql(self, sql):
        result = await self.session.execute(text(sql))
        return [dict(row) for row in result.mappings().fetchall()]

7.3 column_qdrant_repository

column_qdrant_repository用于实现qdrant中的column_info集合的读写操作。

python
from dataclasses import asdict
from qdrant_client import AsyncQdrantClient
from qdrant_client.models import VectorParams, Distance, PointStruct
from app.conf.app_config import app_config
from app.entities.column_info import ColumnInfo


class ColumnQdrantRepository:
    collection_name: str = 'data-agent-column'

    def __init__(self, client: AsyncQdrantClient):
        self.client = client

    async def ensure_collection(self):
        if not await self.client.collection_exists(self.collection_name):
            await self.client.create_collection(self.collection_name,
                vectors_config=VectorParams(size=app_config.qdrant.embedding_size, distance=Distance.COSINE))

    async def upsert(self, ids: list[str], embeddings: list[list[float]], payloads: list[ColumnInfo], batch_size: int = 20):
        zipped = list(zip(ids, embeddings, payloads))
        for i in range(0, len(zipped), batch_size):
            batch = zipped[i:i + batch_size]
            batch_points = [PointStruct(id=id, vector=embedding, payload=asdict(payload))
                           for id, embedding, payload in batch]
            await self.client.upsert(collection_name=self.collection_name, points=batch_points)

    async def search(self, embedding: list[float], score_threshold: float = 0.6, limit: int = 5) -> list[ColumnInfo]:
        result = await self.client.query_points(
            collection_name=self.collection_name, query=embedding,
            score_threshold=score_threshold, limit=limit)
        return [ColumnInfo(**point.payload) for point in result.points]

7.4 metric_qdrant_repository

metric_qdrant_repository用于实现qdrant中的metric_info集合的读写操作。

python
from dataclasses import asdict
from qdrant_client import AsyncQdrantClient
from qdrant_client.models import VectorParams, Distance, PointStruct
from app.conf.app_config import app_config
from app.entities.metric_info import MetricInfo


class MetricQdrantRepository:
    collection_name = 'data-agent-metric'

    def __init__(self, client: AsyncQdrantClient):
        self.client = client

    async def ensure_collection(self):
        if not await self.client.collection_exists(self.collection_name):
            await self.client.create_collection(self.collection_name,
                vectors_config=VectorParams(size=app_config.qdrant.embedding_size, distance=Distance.COSINE))

    async def upsert(self, ids, embeddings, payloads, batch_size=20):
        zipped = list(zip(ids, embeddings, payloads))
        for i in range(0, len(zipped), batch_size):
            batch = zipped[i:i + batch_size]
            batch_points = [PointStruct(id=id, vector=embedding, payload=asdict(payload))
                           for id, embedding, payload in batch]
            await self.client.upsert(collection_name=self.collection_name, points=batch_points)

    async def search(self, embedding, score_threshold=0.6, limit=5) -> list[MetricInfo]:
        result = await self.client.query_points(
            collection_name=self.collection_name, query=embedding,
            score_threshold=score_threshold, limit=limit)
        return [MetricInfo(**point.payload) for point in result.points]

7.5 value_es_repository

value_es_repository用于实现es中的value_info index的读写操作。

python
from dataclasses import asdict
from elasticsearch import AsyncElasticsearch
from app.entities.value_info import ValueInfo


class ValueESRepository:
    index_name = 'data-agent-value'
    index_mappings = {
        "dynamic": False,
        "properties": {
            "id": {"type": "keyword"},
            "value": {"type": "text", "analyzer": "ik_max_word", "search_analyzer": "ik_max_word"},
            "column_id": {"type": "keyword"}
        }
    }

    def __init__(self, client: AsyncElasticsearch):
        self.client = client

    async def ensure_index(self):
        if not await self.client.indices.exists(index=self.index_name):
            await self.client.indices.create(index=self.index_name, mappings=self.index_mappings)

    async def index(self, value_infos: list[ValueInfo], batch_size=20):
        for i in range(0, len(value_infos), batch_size):
            batch = value_infos[i:i + batch_size]
            operations = []
            for value_info in batch:
                operations.append({"index": {"_index": self.index_name, "_id": value_info.id}})
                operations.append(asdict(value_info))
            await self.client.bulk(operations=operations)

    async def search(self, keyword: str, score_threshold: float = 0.6, limit: int = 5) -> list[ValueInfo]:
        result = await self.client.search(index=self.index_name,
            query={"match": {"value": keyword}},
            min_score=score_threshold, size=limit)
        return [ValueInfo(**hit['_source']) for hit in result['hits']['hits']]

八、测试

在终端执行如下命令执行build脚本:

bash
PS D:\code\workspace-python\data-agent> python -m app.scripts.build_meta_knowledge --conf D:\code\workspace-python\data-agent\conf\meta_config.yaml

或者在idea开发工具中执行:tips 需要要设置配置文件路径

image-20260225142446612

执行完毕后,可查看meta数据库、es和qdrant中是否有数据写入。


九、企业痛点-方案映射

痛点传统方案AI Agent 方案效率提升
数仓表结构散落在不同系统人工整理文档自动扫描 dw 库 + 写入 meta 库文档生成从天→分钟
字段语义理解困难逐字段阅读注释为字段名+别名+描述建向量索引语义检索秒级返回
维度值不明确查数仓字典维度取值自动同步到 ES 全文索引查询准确率提升 3x

十、本阶段文件索引

优先级文件路径
🔥 P0build_meta_knowledge.pyapp/scripts/build_meta_knowledge.py
🔥 P0meta_knowledge_service.pyapp/services/meta_knowledge_service.py
🔥 P0meta_config.yamlconf/meta_config.yaml
🟡 P1entities/*.py (5 个)app/entities/
🟡 P1models/*.py (5 个)app/models/
🟡 P1mappers/*.py (4 个)app/repositories/mysql/meta/mappers/
🟡 P1meta_mysql_repository.pyapp/repositories/mysql/meta/meta_mysql_repository.py
🟡 P1dw_mysql_repository.pyapp/repositories/mysql/dw/dw_mysql_repository.py
🟡 P1column_qdrant_repository.pyapp/repositories/qdrant/column_qdrant_repository.py
🟡 P1metric_qdrant_repository.pyapp/repositories/qdrant/metric_qdrant_repository.py
🟡 P1value_es_repository.pyapp/repositories/es/value_es_repository.py
🟢 P2meta_config.pyapp/conf/meta_config.py

OPC 超级个体实战指南