商城首页欢迎来到中国正版软件门户

您的位置:首页 >08 | 把维度值同步到 Elasticsearch(生成阶段)

08 | 把维度值同步到 Elasticsearch(生成阶段)

  发布于2026-07-30 阅读(0)

扫一扫,手机访问

08 | 把维度值同步到 Elasticsearch(生成阶段)

本文是系列文的第八篇,建议按顺序阅读。

08 | 把维度值同步到 Elasticsearch(生成阶段)

本文目标

上一篇已经完成了字段和指标的向量化:把 name + description + alias 转换成向量并写入 Qdrant,用于召回相关字段和指标。

这一篇继续完成生成阶段的第 3 步。举个例子,用户问:

 复制代码北京今年的销售额是多少?

生成 SQL 时,不仅要知道“销售额”对应 order_amount,还得知道“北京”是 dim_region.province 中真实存在的值。换句话说,需要把维度值也同步到某种存储里,方便后续精确匹配:

 复制代码“销售额” → Qdrant → fact_order.order_amount
“北京”   → Elasticsearch → dim_region.province = 北京市

本文正文只讲本项目怎样同步维度值。Elasticsearch 的通用概念统一放在文末的科普模块,方便新手快速了解。

1. 先看完整同步链路

 复制代码meta_config.yaml
└── 找到 sync: true 的字段
        │
        ▼
MySQL DW
└── 查询低基数字段的全部 DISTINCT 值
        │
        ▼
组装维度值文档
        │
        ▼
DimValueSyncService
├── 重建 ES Index
├── 按字段 Bulk 写入
└── 最后统一 Refresh
        │
        ▼
Elasticsearch dim_value

各层职责如下:

文件职责
Infrastructureapp/infrastructure/es_client.py创建 ES 官方异步客户端
Repositoryapp/repositories/dw_db_repository.py读取 DW 低基数字段的全部去重值
Serviceapp/services/dw_db_service.py对外提供 DW 维度值读取能力
Serviceapp/services/dim_value_sync_service.py重建索引、Bulk 写入、刷新索引
Scriptconf/sync_db.py筛选字段、组装文档、装配依赖、释放资源

2. 一条维度值对应一个文档

实际项目中,采用“一条不同维度值对应一个 ES 文档”的方式:

 复制代码{
  "id": "dim_region.province.北京市",
  "column_id": "dim_region.province",
  "value": "北京市"
}

为什么不用数组方案?比如下面这种:

 复制代码{
  "column_id": "dim_region.province",
  "values": ["北京市", "上海市", "广东省"]
}

数组虽然也能搜索,但存在几个比较麻烦的问题:

  • 搜索命中后不容易判断数组中具体匹配了哪个值
  • 高基数字段的数组可能变成非常大的单文档
  • 增删一个值需要更新整个字段文档
  • _id=column_id 只能表示字段,不能稳定表示每个维度值

独立文档更适合后续返回,结构也更清晰:

 复制代码id = dim_region.province.北京市
column_id = dim_region.province
value = 北京市

其中 id 是方便日志和下游识别的可读业务 ID;而真正写入 ES 元数据的 _id 使用固定长度 UUID。

3. 配置 Elasticsearch

先安装官方异步客户端:

 复制代码uv add "elasticsearch[async]>=8,<9"

然后在 conf/app_config.yaml 中配置:

 复制代码es:
  host: localhost
  port: 9200
  index_name: dim_value
  timeout: 60

各参数作用:

参数作用
host / portES HTTP 服务地址
index_name保存维度值的索引名
timeout单次 ES HTTP 请求超时秒数

4. 定义 ESClient

文件 app/infrastructure/es_client.py

 复制代码from elasticsearch import AsyncElasticsearch
class ESClient(AsyncElasticsearch):
    """统一项目构造参数的 ES 官方异步客户端薄适配层。"""    def __init__(self, base_url: str, timeout: float = 60.0):
        super().__init__(
            hosts=base_url.rstrip("/"),
            request_timeout=timeout,
        )

注意,Elasticsearch 8 使用 request_timeout 设置单次请求超时,不要沿用旧示例里的 timeout=int(timeout) 写法。

AsyncElasticsearch 内部管理 HTTP 连接池,在一次同步任务中复用,任务结束后统一关闭:

 复制代码await es_client.close()

目前 ES 的逻辑只有重建、写入和刷新,所以 DimValueSyncService 直接调用 ES Client 即可,暂时不增加 ES Repository。以后查询、过滤、聚合等数据访问操作明显增多时,再抽取 Repository 也不迟。

5. 定义 DimValueSyncService

文件:app/services/dim_value_sync_service.py

重建索引并定义 Mapping

 复制代码async def reset_index(self) -> None:
    await self.es_client.indices.delete(
        index=self.index_name,
        ignore_unavailable=True,
    )
    await self.es_client.indices.create(
        index=self.index_name,
        settings={
            "number_of_shards": 1,
            "number_of_replicas": 0,
        },
        mappings={
            "properties": {
                "id": {"type": "keyword"},
                "column_id": {"type": "keyword"},
                "value": {
                    "type": "text",
                    "analyzer": "ik_max_word",
                    "search_analyzer": "ik_smart",
                    "fields": {
                        "keyword": {
                            "type": "keyword",
                            "ignore_above": 512,
                        }
                    },
                },
            }
        },
    )

字段的用途:

字段类型作用
idkeyword可读业务 ID,格式为 column_id.value
column_idkeyword精确定位 表名.字段名
valuetext + keyword同时支持中文全文检索和精确匹配

table_idcolumn_name 都能从 column_id 推导;字段描述和别名已经由 Qdrant 负责语义召回,所以 ES 里不需要重复保存。

当前 Docker 是单节点,所以副本数设置为 0。否则副本分片无法分配,集群会长期显示黄色,看着揪心。

而且这是生成阶段的全量重建策略:每次删除旧索引再创建新索引,保证已经从 DW 删除的旧值不会残留。

使用原生 Bulk API

原生 Bulk API 每条文档需要两个对象:

 复制代码操作元数据
文档内容

项目中的实现:

 复制代码async def index_values(
    self,
    documents: list[dict[str, Any]],
) -> int:
    if not documents:
        return 0    operations = []
    for document in documents:
        document_id = self.stable_document_id(
            document["column_id"],
            document["value"],
        )
        operations.extend([
            {
                "index": {
                    "_index": self.index_name,
                    "_id": document_id,
                }
            },
            document,
        ])    response = await self.es_client.bulk(operations=operations)
    if response.get("errors"):
        raise RuntimeError("Elasticsearch Bulk 写入失败")
    return len(documents)

原来的代码混用了两种 Bulk 格式:

 复制代码{
    "index": {
        "_op_type": "index",
        "_source": document,
    }
}

_op_type_source 属于 helpers.async_bulk() 的 Action 格式;直接调用 client.bulk() 时必须使用上面的 Operations 格式,不能混用。

检查部分失败

HTTP 请求成功不代表每个文档都成功。Bulk 响应可能出现部分失败:

 复制代码{
  "errors": true,
  "items": []
}

因此项目会检查 response["errors"],发现部分失败时抛出异常,并展示首个错误,避免把“部分成功”误认为“全部成功”。

稳定文档 ID

同一字段有多个值,所以 _id 不能只使用 column_id。项目根据 column_id + value 生成稳定 UUID:

 复制代码@staticmethod
def stable_document_id(column_id: str, value: str) -> str:
    return str(uuid5(NAMESPACE_URL, f"{column_id}:{value}"))

相同字段和值每次得到相同 ID,重复同步不会产生重复文档。

写入后统一 Refresh

Bulk 默认是近实时可搜索,不需要每批都刷新。全部写完后统一执行一次:

 复制代码async def refresh(self) -> None:
    await self.es_client.indices.refresh(index=self.index_name)

这比每一批都 Refresh 更高效。

6. 读取 DW 维度值

下面的代码只是限制读取数量,并不是真正的全量读取:

 复制代码await get_column_examples(..., limit=100000)

本项目明确只允许低基数字段设置 sync: true,每个字段的不同值数量很少,因此不需要增加分页循环。

DwDBRepository 直接读取该字段的全部非空去重值:

 复制代码async def get_distinct_column_values(
    self,
    table_name: str,
    column_name: str,
) -> list[str]:
    self._validate_identifier(table_name)
    self._validate_identifier(column_name)    sql = text(
        f"SELECT DISTINCT `{column_name}` "
        f"FROM `{table_name}` "
        f"WHERE `{column_name}` IS NOT NULL "
        f"ORDER BY `{column_name}`"
    )
    result = await self.session.execute(sql)
    return [str(row[0]) for row in result.all()]

这里有几个细节需要注意:

  • DISTINCT:同一个值只同步一次
  • IS NOT NULL:不写入空值
  • ORDER BY:保证返回顺序稳定
  • 表名和字段名不能参数化,必须先校验再拼接

DwDBService 提供同名方法,把 Repository 能力暴露给同步流程。

7. 编写同步脚本

文件:conf/sync_db.py

 复制代码async def sync_to_elasticsearch(
    config: dict,
    dw_database: MySQLDatabase,
) -> int:
    es_config = app_config["es"]    es_client = ESClient(
        base_url=f"http://{es_config['host']}:{es_config['port']}",
        timeout=es_config.get("timeout", 60),
    )
    service = DimValueSyncService(
        es_client=es_client,
        index_name=es_config["index_name"],
    )
    total_count = 0    try:
        await service.reset_index()        async with dw_database.session() as dw_session:
            dw_service = DwDBService(DwDBRepository(dw_session))            for table in config.get("tables", []):
                table_name = table["name"]
                for column in table.get("columns", []):
                    if not column.get("sync", False):
                        continue                    column_name = column["name"]
                    column_id = f"{table_name}.{column_name}"
                    values = await dw_service.get_distinct_column_values(
                        table_name,
                        column_name,
                    )                    documents = [
                        {
                            "id": f"{column_id}.{value}",
                            "column_id": column_id,
                            "value": value,
                        }
                        for value in values
                        if value.strip()
                    ]
                    total_count += await service.index_values(documents)        await service.refresh()
        return total_count
    finally:
        await es_client.close()

这里修正了原实现中的几个问题:

  • 删除了没有使用的 session 参数
  • 返回类型改为实际的 int
  • 只处理 sync: true 的字段,不为其他字段构造空文档
  • 使用现有的 host/port/index_name 配置
  • 每个低基数字段使用一次 Bulk 请求写入,不逐值发送 HTTP 请求
  • finally 中关闭 ES Client

主流程必须使用 await

 复制代码dim_value_count = await sync_to_elasticsearch(
    meta_config,
    dw_database,
)
print(f"已写入 Elasticsearch 维度值 {dim_value_count} 条")

如果漏掉 await,函数只会返回一个协程对象,同步逻辑不会真正执行。

8. 校验 sync 配置

sync 必须是布尔值,而且只有维度字段才适合同步枚举值:

 复制代码sync_enabled = column.get("sync", False)
if not isinstance(sync_enabled, bool):
    raise ValueError("sync 必须是布尔值")
if sync_enabled and column["role"] != "dimension":
    raise ValueError("只有 dimension 字段才能设置 sync=true")

这样可以避免把字符串 "false" 当成真值,也能避免意外把订单金额、主键等高基数字段全部写入 ES。

9. 执行与验证

确认 Elasticsearch 健康:

 复制代码curl 

执行完整生成阶段同步:

 复制代码uv run python conf/sync_db.py

查看文档数量:

 复制代码curl 

搜索“北京”:

 复制代码curl -X POST  
  -H 'Content-Type: application/json' 
  -d '{
    "query": {"match": {"value": "北京"}},
    "_source": ["column_id", "value"]
  }'

本项目实际验证结果:

 复制代码同步返回数量:75
ES 索引数量:75
“北京”召回:dim_region.province → 北京市

使用 Kibana 网页查看

Docker Compose 已经启动 Kibana,并映射到本机 5601 端口:

 复制代码

Kibana 不单独保存维度值,它只是用网页方式查看和操作 Elasticsearch。

方法一:在 Dev Tools 中查询

进入 Kibana 的 Dev Tools → Console,查看文档数量:

 复制代码GET dim_value/_count

查看少量维度值:

 复制代码GET dim_value/_search
{
  "size": 10,
  "_source": ["id", "column_id", "value"],
  "query": {
    "match_all": {}
  }
}

搜索“北京”:

 复制代码GET dim_value/_search
{
  "size": 10,
  "_source": ["id", "column_id", "value"],
  "query": {
    "match": {
      "value": "北京"
    }
  }
}
方法二:在 Discover 中浏览

第一次使用时,需要创建 Data View:

  1. 打开 Stack Management → Data Views
  2. 点击 Create data view
  3. Name 填写 dim_value
  4. Index pattern 填写 dim_value
  5. 本索引没有时间字段,时间字段选择“不使用时间过滤器”。
  6. 创建后进入 Discover,选择 dim_value Data View。

在 Discover 中可以直接查看 idcolumn_idvalue,也可以在搜索框中输入:

 复制代码value: 北京

如果 暂时打不开,可以先检查:

 复制代码docker ps --filter name=kibana
docker logs --tail 100 kibana

10. 项目内常见问题

现象原因处理
连接 9200 失败ES 未启动或端口错误检查 Docker、端口和集群健康状态
创建索引提示找不到 ik_max_word没有安装 IK 插件使用项目提供的 ES 镜像或安装匹配版本的 IK
写入后立刻搜不到ES 是近实时搜索全部写完后调用一次 indices.refresh()
Bulk HTTP 成功但少了文档部分 Action 失败检查 Bulk 响应中的 errorsitems
重复同步产生重复文档_id 不稳定根据 column_id + value 生成稳定 ID
已删除的旧值仍存在只执行 Upsert开发阶段全量重建索引
内存占用过高sync: true 字段不再是低基数恢复分页或改用游标、流式 Bulk
ES 集群显示黄色单节点却配置了副本开发单节点设置 number_of_replicas: 0

科普 Elasticsearch

本模块只讲通用概念,不依赖当前项目代码。

1. 什么是 Elasticsearch

Elasticsearch(简称 ES)是基于 Apache Lucene 构建的分布式搜索和分析引擎。应用通过 HTTP API 写入 JSON 文档,然后可以对海量文档执行:

  • 全文检索
  • 精确过滤
  • 相关性排序
  • 聚合统计
  • 高亮显示
  • 自动补全

典型搜索过程:

 复制代码应用写入 JSON 文档
→ Elasticsearch 分词并建立索引
→ 用户提交查询
→ Elasticsearch 找到候选文档并计算相关性
→ 返回排序后的结果

Elasticsearch 的核心目标是“搜索和分析”,不是处理复杂关系、强事务和多表关联。很多系统会同时使用关系数据库和 Elasticsearch:关系数据库保存权威数据,Elasticsearch 保存适合检索的副本。

2. 核心组成

Cluster 和 Node

  • Cluster(集群):一组共同工作的 Elasticsearch 节点
  • Node(节点):一个正在运行的 Elasticsearch 实例
 复制代码Cluster
├── Node A
├── Node B
└── Node C

开发环境可以只有一个节点,生产环境通常使用多个节点来分担存储、查询和故障恢复。

Index、Document 和 Field

 复制代码Index: articles
└── Document
    ├── title
    ├── content
    ├── author
    └── published_at
Elasticsearch粗略类比关系数据库含义
IndexTable一类文档的逻辑集合
DocumentRow一条 JSON 数据
FieldColumn文档中的一个属性
MappingSchema字段类型和索引规则
Document _idPrimary Key文档唯一标识

这种类比只用于入门。ES 的 Index 还包含倒排索引、分片、Segment 等搜索结构,并不等于关系表。

一条文档示例:

 复制代码{
  "title": "Elasticsearch 入门",
  "content": "介绍全文检索和倒排索引",
  "author": "张三",
  "published_at": "2026-07-30"
}

3. 什么是倒排索引

普通阅读方向是“从文档找到词”:

 复制代码文档 A → 北京市
文档 B → 上海市

倒排索引反过来保存“从词找到文档”:

 复制代码北京 → 文档 A
上海 → 文档 B

搜索时不需要逐条扫描所有文档,因此大量文本也能快速查询。

4. Mapping 和字段类型

Mapping 定义文档字段怎样被存储和索引,类似数据库中的 Schema,但更关注搜索行为:

 复制代码{
  "mappings": {
    "properties": {
      "title": {"type": "text"},
      "author": {"type": "keyword"},
      "price": {"type": "double"},
      "published_at": {"type": "date"}
    }
  }
}

常见字段类型:

类型常见用途
text全文检索文本
keyword精确值、过滤、排序和聚合
integer / long整数
float / double小数
boolean布尔值
date日期和时间
object普通嵌套 JSON 对象
nested需要保持对象独立关系的对象数组

ES 可以通过 Dynamic Mapping 自动推断字段类型,但生产系统最好为重要索引显式定义 Mapping。自动推断错误后,已有字段类型通常不能直接修改,需要新建索引并 Reindex。

5. Analyzer、Tokenizer 和 Token

Analyzer 决定文本写入和查询时怎样被拆成 Token,通常包含:

 复制代码Character Filter
→ Tokenizer
→ Token Filter

例如:

 复制代码“北京市朝阳区”
→ 北京、北京市、朝阳、朝阳区……

中文不能简单按空格分词,常见方案之一是安装 IK 分词插件:

  • ik_max_word:写入时尽可能细分,增加可召回词
  • ik_smart:查询时做较粗粒度切分,降低噪声

分词器必须在创建 Mapping 时确定。修改 Analyzer 后通常需要重建索引。

IK 属于第三方插件,插件版本通常必须与 Elasticsearch 版本一致。例如 Elasticsearch 8.19.10 应搭配对应的 IK 8.19.10 插件。没有安装插件却在 Mapping 中使用 ik_max_word,创建索引时会失败。

可以通过 Analyze API 查看分词结果:

 复制代码POST _analyze
{
  "analyzer": "ik_max_word",
  "text": "北京市朝阳区"
}

6. text 和 keyword

这是 ES 最重要的字段类型区别之一:

类型是否分词常见用途
text全文搜索、相关性匹配
keyword精确过滤、排序、聚合

同一个字段可以同时支持两种能力:

 复制代码{
  "title": {
    "type": "text",
    "fields": {
      "keyword": {"type": "keyword"}
    }
  }
}

查询方式:

 复制代码match title=搜索引擎              → 全文检索
term title.keyword=搜索引擎入门    → 精确匹配

不要对 text 字段直接使用 term 查询,因为它保存的是分词后的 Token。

7. Query DSL:Match、Term 和 Bool

全文检索:

 复制代码{
  "query": {
    "match": {
      "title": "搜索引擎"
    }
  }
}

精确过滤:

 复制代码{
  "query": {
    "term": {
      "status": "published"
    }
  }
}

组合查询:

 复制代码{
  "query": {
    "bool": {
      "must": [
        {"match": {"title": "搜索引擎"}}
      ],
      "filter": [
        {"term": {"status": "published"}},
        {"range": {"published_at": {"gte": "2026-01-01"}}}
      ]
    }
  }
}

must 参与相关性评分;filter 只判断是否满足条件,适合结构化约束。

其他常用查询:

  • multi_match:同时搜索多个文本字段
  • range:数字或日期范围
  • exists:字段是否存在
  • prefix:前缀匹配
  • wildcard:通配符匹配,应谨慎使用

8. 相关性评分

全文检索结果通常带有 _score。ES 默认使用 BM25 等算法,根据以下因素评估文档相关性:

  • 查询词是否出现
  • 查询词在文档中出现多少次
  • 查询词在整个索引中是否稀有
  • 字段文本长度
  • 查询中的 Boost 权重
 复制代码更相关的文档 → 更高的 _score → 排在更前面

_score 不是正确概率,不能把 1.8 理解成 180% 或把 0.8 理解成 80% 正确。评分只适合在同一次查询结果中比较排序。

9. Bulk API

逐条发送 HTTP 请求会产生大量网络往返。Bulk API 把多条操作放进一个请求:

 复制代码index 元数据
文档 1
index 元数据
文档 2

两种 Python 写法不要混淆:

  • client.bulk(operations=...):使用原生 Operations 格式
  • helpers.async_bulk(client, actions):使用带 _op_type_source 的 Action 格式

它们都能批量写入,但请求数据结构不同。

批次不是越大越好,需要综合考虑:

  • 单条文档大小
  • HTTP 请求体大小
  • ES Heap
  • 网络延迟
  • Bulk 响应时间

常见起点是每批 500~1000 条,再根据真实数据测试。

10. 近实时搜索、Refresh 和 Segment

ES 写入成功后,文档通常不会在同一瞬间就能被搜索,这叫 Near Real-Time(近实时)。

 复制代码写入成功
→ 等待 Refresh
→ 新 Segment 可搜索

频繁手动 Refresh 会降低写入吞吐。批量任务更适合全部写完后刷新一次;在线写入通常依赖 ES 的自动刷新周期。

Lucene 会把可搜索数据组织成不可变的 Segment。新数据经过 Refresh 后形成可搜索 Segment,后台还会自动合并小 Segment。理解这一点有助于解释:

  • 为什么 ES 搜索是近实时而不是实时
  • 为什么频繁 Refresh 会降低写入性能
  • 为什么大量更新和删除不会立刻释放磁盘空间

11. Shard 和 Replica

  • Primary Shard:索引数据的主分片
  • Replica Shard:主分片的副本,用于容灾和扩展查询吞吐

单节点开发环境无法把副本放到另一台节点,所以通常设置:

 复制代码{
  "number_of_shards": 1,
  "number_of_replicas": 0
}

生产集群有多个节点时,再根据数据量、吞吐和容灾要求调整。

12. Aggregation 聚合

Aggregation 用来统计和分析文档,类似 SQL 中的 GROUP BYCOUNTAVG 等操作。

统计不同状态的文档数量:

 复制代码{
  "size": 0,
  "aggs": {
    "status_count": {
      "terms": {
        "field": "status"
      }
    }
  }
}

常用聚合:

  • terms:按精确值分组
  • date_histogram:按时间区间分组
  • range:按范围分组
  • avg / sum / min / max:数值统计
  • cardinality:估算不同值数量

参与分组、排序和聚合的字符串字段通常应该使用 keyword,而不是经过分词的 text

13. 索引生命周期、Reindex 和 Alias

Mapping 或 Analyzer 发生变化时,通常不能直接修改已有数据结构。常见迁移流程:

 复制代码创建 articles_v2
→ 写入新数据或执行 Reindex
→ 验证数量和查询结果
→ Alias 从 articles_v1 切换到 articles_v2
→ 删除旧索引

Alias 是指向一个或多个真实索引的逻辑名称。应用始终查询 Alias,可以在不修改应用配置的情况下切换索引版本。

对于持续增长的日志、指标等时序数据,还可以使用 Index Lifecycle Management(ILM)自动完成滚动、迁移和删除。

14. Elasticsearch 适合与不适合什么

适合的场景:

  • 网站、商品和文章全文搜索
  • 日志检索与可观测性分析
  • 多条件过滤和相关性排序
  • 实时或近实时聚合分析
  • 自动补全、高亮和同义词搜索

不适合作为首选方案的场景:

  • 强事务和严格 ACID 一致性
  • 大量跨实体 Join
  • 频繁的小粒度随机更新
  • 只按主键进行简单读写
  • 把 ES 当作唯一且不可恢复的数据源

选择技术时应先看查询模式。需要事务和关系约束时优先考虑关系数据库;需要文本检索和搜索分析时再考虑 Elasticsearch。

15. 常用管理操作

下面的命令可以在 Kibana Dev Tools 中执行:

 复制代码# 查看集群健康状态
GET _cluster/health# 查看索引列表
GET _cat/indices?v# 查看 Mapping
GET articles/_mapping# 查看文档数量
GET articles/_count# 查询少量文档
GET articles/_search
{
  "size": 10,
  "query": {"match_all": {}}
}

学习 Elasticsearch 时,建议按以下顺序掌握:

 复制代码Document 与 Mapping
→ text / keyword
→ Analyzer 与倒排索引
→ Match / Term / Bool
→ Bulk 与 Refresh
→ Shard / Replica
→ Aggregation
→ Alias 与索引迁移
本文转载于:https://juejin.cn/post/7668148348428615706 如有侵犯,请联系zhengruancom@outlook.com删除。
免责声明:正软商城发布此文仅为传递信息,不代表正软商城认同其观点或证实其描述。

热门关注