您的位置:首页 >08 | 把维度值同步到 Elasticsearch(生成阶段)
发布于2026-07-30 阅读(0)
扫一扫,手机访问
本文是系列文的第八篇,建议按顺序阅读。

上一篇已经完成了字段和指标的向量化:把 name + description + alias 转换成向量并写入 Qdrant,用于召回相关字段和指标。
这一篇继续完成生成阶段的第 3 步。举个例子,用户问:
复制代码北京今年的销售额是多少?
生成 SQL 时,不仅要知道“销售额”对应 order_amount,还得知道“北京”是 dim_region.province 中真实存在的值。换句话说,需要把维度值也同步到某种存储里,方便后续精确匹配:
复制代码“销售额” → Qdrant → fact_order.order_amount
“北京” → Elasticsearch → dim_region.province = 北京市
本文正文只讲本项目怎样同步维度值。Elasticsearch 的通用概念统一放在文末的科普模块,方便新手快速了解。
复制代码meta_config.yaml
└── 找到 sync: true 的字段
│
▼
MySQL DW
└── 查询低基数字段的全部 DISTINCT 值
│
▼
组装维度值文档
│
▼
DimValueSyncService
├── 重建 ES Index
├── 按字段 Bulk 写入
└── 最后统一 Refresh
│
▼
Elasticsearch dim_value
各层职责如下:
| 层 | 文件 | 职责 |
|---|---|---|
| Infrastructure | app/infrastructure/es_client.py | 创建 ES 官方异步客户端 |
| Repository | app/repositories/dw_db_repository.py | 读取 DW 低基数字段的全部去重值 |
| Service | app/services/dw_db_service.py | 对外提供 DW 维度值读取能力 |
| Service | app/services/dim_value_sync_service.py | 重建索引、Bulk 写入、刷新索引 |
| Script | conf/sync_db.py | 筛选字段、组装文档、装配依赖、释放资源 |
实际项目中,采用“一条不同维度值对应一个 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。
先安装官方异步客户端:
复制代码uv add "elasticsearch[async]>=8,<9"
然后在 conf/app_config.yaml 中配置:
复制代码es:
host: localhost
port: 9200
index_name: dim_value
timeout: 60
各参数作用:
| 参数 | 作用 |
|---|---|
host / port | ES HTTP 服务地址 |
index_name | 保存维度值的索引名 |
timeout | 单次 ES HTTP 请求超时秒数 |
文件 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 也不迟。
文件:app/services/dim_value_sync_service.py
复制代码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,
}
},
},
}
},
)
字段的用途:
| 字段 | 类型 | 作用 |
|---|---|---|
id | keyword | 可读业务 ID,格式为 column_id.value |
column_id | keyword | 精确定位 表名.字段名 |
value | text + keyword | 同时支持中文全文检索和精确匹配 |
table_id 和 column_name 都能从 column_id 推导;字段描述和别名已经由 Qdrant 负责语义召回,所以 ES 里不需要重复保存。
当前 Docker 是单节点,所以副本数设置为 0。否则副本分片无法分配,集群会长期显示黄色,看着揪心。
而且这是生成阶段的全量重建策略:每次删除旧索引再创建新索引,保证已经从 DW 删除的旧值不会残留。
原生 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 不能只使用 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,重复同步不会产生重复文档。
Bulk 默认是近实时可搜索,不需要每批都刷新。全部写完后统一执行一次:
复制代码async def refresh(self) -> None:
await self.es_client.indices.refresh(index=self.index_name)
这比每一批都 Refresh 更高效。
下面的代码只是限制读取数量,并不是真正的全量读取:
复制代码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 能力暴露给同步流程。
文件: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 参数intsync: true 的字段,不为其他字段构造空文档host/port/index_name 配置finally 中关闭 ES Client主流程必须使用 await:
复制代码dim_value_count = await sync_to_elasticsearch(
meta_config,
dw_database,
)
print(f"已写入 Elasticsearch 维度值 {dim_value_count} 条")
如果漏掉 await,函数只会返回一个协程对象,同步逻辑不会真正执行。
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。
确认 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 → 北京市
Docker Compose 已经启动 Kibana,并映射到本机 5601 端口:
复制代码
Kibana 不单独保存维度值,它只是用网页方式查看和操作 Elasticsearch。
进入 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": "北京"
}
}
}
第一次使用时,需要创建 Data View:
dim_value。dim_value。dim_value Data View。在 Discover 中可以直接查看 id、column_id、value,也可以在搜索框中输入:
复制代码value: 北京
如果 暂时打不开,可以先检查:
复制代码docker ps --filter name=kibana
docker logs --tail 100 kibana
| 现象 | 原因 | 处理 |
|---|---|---|
连接 9200 失败 | ES 未启动或端口错误 | 检查 Docker、端口和集群健康状态 |
创建索引提示找不到 ik_max_word | 没有安装 IK 插件 | 使用项目提供的 ES 镜像或安装匹配版本的 IK |
| 写入后立刻搜不到 | ES 是近实时搜索 | 全部写完后调用一次 indices.refresh() |
| Bulk HTTP 成功但少了文档 | 部分 Action 失败 | 检查 Bulk 响应中的 errors 和 items |
| 重复同步产生重复文档 | _id 不稳定 | 根据 column_id + value 生成稳定 ID |
| 已删除的旧值仍存在 | 只执行 Upsert | 开发阶段全量重建索引 |
| 内存占用过高 | sync: true 字段不再是低基数 | 恢复分页或改用游标、流式 Bulk |
| ES 集群显示黄色 | 单节点却配置了副本 | 开发单节点设置 number_of_replicas: 0 |
本模块只讲通用概念,不依赖当前项目代码。
Elasticsearch(简称 ES)是基于 Apache Lucene 构建的分布式搜索和分析引擎。应用通过 HTTP API 写入 JSON 文档,然后可以对海量文档执行:
典型搜索过程:
复制代码应用写入 JSON 文档
→ Elasticsearch 分词并建立索引
→ 用户提交查询
→ Elasticsearch 找到候选文档并计算相关性
→ 返回排序后的结果
Elasticsearch 的核心目标是“搜索和分析”,不是处理复杂关系、强事务和多表关联。很多系统会同时使用关系数据库和 Elasticsearch:关系数据库保存权威数据,Elasticsearch 保存适合检索的副本。
复制代码Cluster
├── Node A
├── Node B
└── Node C
开发环境可以只有一个节点,生产环境通常使用多个节点来分担存储、查询和故障恢复。
复制代码Index: articles
└── Document
├── title
├── content
├── author
└── published_at
| Elasticsearch | 粗略类比关系数据库 | 含义 |
|---|---|---|
| Index | Table | 一类文档的逻辑集合 |
| Document | Row | 一条 JSON 数据 |
| Field | Column | 文档中的一个属性 |
| Mapping | Schema | 字段类型和索引规则 |
Document _id | Primary Key | 文档唯一标识 |
这种类比只用于入门。ES 的 Index 还包含倒排索引、分片、Segment 等搜索结构,并不等于关系表。
一条文档示例:
复制代码{
"title": "Elasticsearch 入门",
"content": "介绍全文检索和倒排索引",
"author": "张三",
"published_at": "2026-07-30"
}
普通阅读方向是“从文档找到词”:
复制代码文档 A → 北京市
文档 B → 上海市
倒排索引反过来保存“从词找到文档”:
复制代码北京 → 文档 A
上海 → 文档 B
搜索时不需要逐条扫描所有文档,因此大量文本也能快速查询。
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。
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": "北京市朝阳区"
}
这是 ES 最重要的字段类型区别之一:
| 类型 | 是否分词 | 常见用途 |
|---|---|---|
text | 是 | 全文搜索、相关性匹配 |
keyword | 否 | 精确过滤、排序、聚合 |
同一个字段可以同时支持两种能力:
复制代码{
"title": {
"type": "text",
"fields": {
"keyword": {"type": "keyword"}
}
}
}
查询方式:
复制代码match title=搜索引擎 → 全文检索
term title.keyword=搜索引擎入门 → 精确匹配
不要对 text 字段直接使用 term 查询,因为它保存的是分词后的 Token。
全文检索:
复制代码{
"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:通配符匹配,应谨慎使用全文检索结果通常带有 _score。ES 默认使用 BM25 等算法,根据以下因素评估文档相关性:
复制代码更相关的文档 → 更高的 _score → 排在更前面
_score 不是正确概率,不能把 1.8 理解成 180% 或把 0.8 理解成 80% 正确。评分只适合在同一次查询结果中比较排序。
逐条发送 HTTP 请求会产生大量网络往返。Bulk API 把多条操作放进一个请求:
复制代码index 元数据
文档 1
index 元数据
文档 2
两种 Python 写法不要混淆:
client.bulk(operations=...):使用原生 Operations 格式helpers.async_bulk(client, actions):使用带 _op_type、_source 的 Action 格式它们都能批量写入,但请求数据结构不同。
批次不是越大越好,需要综合考虑:
常见起点是每批 500~1000 条,再根据真实数据测试。
ES 写入成功后,文档通常不会在同一瞬间就能被搜索,这叫 Near Real-Time(近实时)。
复制代码写入成功
→ 等待 Refresh
→ 新 Segment 可搜索
频繁手动 Refresh 会降低写入吞吐。批量任务更适合全部写完后刷新一次;在线写入通常依赖 ES 的自动刷新周期。
Lucene 会把可搜索数据组织成不可变的 Segment。新数据经过 Refresh 后形成可搜索 Segment,后台还会自动合并小 Segment。理解这一点有助于解释:
单节点开发环境无法把副本放到另一台节点,所以通常设置:
复制代码{
"number_of_shards": 1,
"number_of_replicas": 0
}
生产集群有多个节点时,再根据数据量、吞吐和容灾要求调整。
Aggregation 用来统计和分析文档,类似 SQL 中的 GROUP BY、COUNT、AVG 等操作。
统计不同状态的文档数量:
复制代码{
"size": 0,
"aggs": {
"status_count": {
"terms": {
"field": "status"
}
}
}
}
常用聚合:
terms:按精确值分组date_histogram:按时间区间分组range:按范围分组avg / sum / min / max:数值统计cardinality:估算不同值数量参与分组、排序和聚合的字符串字段通常应该使用 keyword,而不是经过分词的 text。
Mapping 或 Analyzer 发生变化时,通常不能直接修改已有数据结构。常见迁移流程:
复制代码创建 articles_v2
→ 写入新数据或执行 Reindex
→ 验证数量和查询结果
→ Alias 从 articles_v1 切换到 articles_v2
→ 删除旧索引
Alias 是指向一个或多个真实索引的逻辑名称。应用始终查询 Alias,可以在不修改应用配置的情况下切换索引版本。
对于持续增长的日志、指标等时序数据,还可以使用 Index Lifecycle Management(ILM)自动完成滚动、迁移和删除。
适合的场景:
不适合作为首选方案的场景:
选择技术时应先看查询模式。需要事务和关系约束时优先考虑关系数据库;需要文本检索和搜索分析时再考虑 Elasticsearch。
下面的命令可以在 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 与索引迁移
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
售后无忧
立即购买>office旗舰店
正版软件
正版软件
正版软件
正版软件
正版软件
1
2
3
7
8