S3-13 搜索引擎设计:6 种搜索模式 + 分面 + 热词建议
智策平台的搜索引擎 SearchService 提供 6 种搜索模式(BESTMATCH / PREFIX / FUZZY / EXACT / WILDCARD / REGEX),支持分面搜索(TERMS / RANGE / DATERANGE)、基于 Redis 的热词与历史建议、保存搜索、权限感知过滤,并深度集成 Doris 倒排索引。本文从架构设计到实现细节,完整拆解搜索引擎的每一层。
S3-13 搜索引擎设计:6 种搜索模式 + 分面 + 热词建议
“系列:S3 数据基座 · 第 13 篇 | 难度:高级 | 阅读时间:20 分钟
#TL;DR
智策平台的搜索引擎 SearchService 提供 6 种搜索模式(BEST_MATCH / PREFIX / FUZZY / EXACT / WILDCARD / REGEX),支持分面搜索(TERMS / RANGE / DATE_RANGE)、基于 Redis 的热词与历史建议、保存搜索、权限感知过滤,并深度集成 Doris 倒排索引。本文从架构设计到实现细节,完整拆解搜索引擎的每一层。
#1. 为什么 Ontology 平台需要自己的搜索引擎
在传统数据平台中,搜索往往被降级为"按名称过滤"。但在 Ontology 驱动的系统中,搜索是核心交互入口:
- 对象发现:用户需要在数百万实体中快速定位目标对象
- 关系探索:搜索结果需要展示对象间的关系上下文
- 元数据检索:搜索不仅是数据搜索,还包含类型定义、属性描述、标签等元数据
- 操作入口:搜索结果直接关联可用的 Action(审批、更新、触发工作流)
Palantir Foundry 的搜索体验是其核心竞争力之一。用户在 Workshop 或 Vertex 中搜索时,结果不仅包含匹配的对象,还包含分面统计、关联关系、可执行动作。智策平台的 SearchService 就是要达到这个级别。
+------------------------------------------------------------------+
| SearchService |
| |
| +-----------+ +-----------+ +-----------+ +----------------+ |
| | Query | | Facet | | Suggest | | Saved Search | |
| | Engine | | Engine | | Engine | | Manager | |
| +-----------+ +-----------+ +-----------+ +----------------+ |
| | | | | |
| +-----------+ +-----------+ +-----------+ +----------------+ |
| | Mode | | Agg | | Redis | | PostgreSQL | |
| | Router | | Builder | | Hot/Hist | | Storage | |
| +-----------+ +-----------+ +-----------+ +----------------+ |
| | | |
| +----+--------------+----------------------------------------+ |
| | Permission-Aware Filter Layer | |
| +------------------------------------------------------------+ |
| | |
| +------------------------------------------------------------+ |
| | Doris Inverted Index / Full-Text Engine | |
| +------------------------------------------------------------+ |
+------------------------------------------------------------------+
#2. 搜索请求模型
#2.1 SearchRequest 定义
message SearchRequest {
string world_id = 1;
string query = 2;
SearchMode mode = 3;
repeated string object_types = 4; // 限定搜索的类型
repeated string search_fields = 5; // 限定搜索的字段
repeated FacetRequest facets = 6;
repeated FilterClause filters = 7;
Pagination pagination = 8;
SortSpec sort = 9;
bool include_highlights = 10;
bool include_relations = 11;
}
enum SearchMode {
BEST_MATCH = 0; // 默认:综合评分
PREFIX = 1; // 前缀匹配
FUZZY = 2; // 模糊匹配(容错)
EXACT = 3; // 精确匹配
WILDCARD = 4; // 通配符
REGEX = 5; // 正则表达式
}
#2.2 SearchResponse 定义
message SearchResponse {
repeated SearchHit hits = 1;
int64 total_count = 2;
repeated FacetResult facets = 3;
repeated Suggestion suggestions = 4;
SearchMetadata metadata = 5;
}
message SearchHit {
string object_type = 1;
string object_id = 2;
double score = 3;
map<string, Value> attributes = 4;
repeated Highlight highlights = 5;
repeated RelatedObject relations = 6;
}
#3. 6 种搜索模式详解
#3.1 BEST_MATCH — 综合评分搜索
BEST_MATCH 是默认模式,也是最复杂的模式。它不是简单的全文匹配,而是综合考虑多个信号进行评分:
Score = w1 * text_relevance
+ w2 * field_boost
+ w3 * recency_score
+ w4 * popularity_score
+ w5 * type_priority
text_relevance 由 BM25 算法计算,Doris 的倒排索引原生支持。field_boost 对不同字段赋予不同权重——标题字段权重 3.0,描述字段权重 1.5,标签字段权重 2.0。recency_score 基于对象最后修改时间衰减。popularity_score 基于对象被查看/引用的次数。
class BestMatchScorer:
"""综合评分器"""
FIELD_WEIGHTS = {
"title": 3.0,
"display_name": 3.0,
"description": 1.5,
"tags": 2.0,
"properties": 1.0,
}
def score(self, query: str, hit: RawHit) -> float:
text_score = hit.bm25_score
field_score = self.FIELD_WEIGHTS.get(hit.matched_field, 1.0)
recency = self._recency_decay(hit.updated_at)
popularity = math.log1p(hit.view_count) * 0.1
type_priority = self._type_priority(hit.object_type)
return (
0.5 * text_score * field_score
+ 0.15 * recency
+ 0.1 * popularity
+ 0.25 * type_priority
)
def _recency_decay(self, updated_at: datetime) -> float:
days_ago = (datetime.utcnow() - updated_at).days
return math.exp(-0.01 * days_ago) # 半衰期 ~70 天
def _type_priority(self, object_type: str) -> float:
# 核心业务类型优先
priorities = {"Customer": 1.0, "Order": 0.9, "Product": 0.85}
return priorities.get(object_type, 0.5)
#3.2 PREFIX — 前缀匹配
前缀搜索是自动补全场景的核心。用户输入 "cust" 就能匹配到 "Customer"、"Custom Field" 等。
-- Doris 前缀查询
SELECT * FROM ontology_objects
WHERE display_name LIKE 'cust%'
ORDER BY length(display_name) ASC, view_count DESC
LIMIT 10;
前缀搜索的关键优化是字段级倒排索引。Doris 支持在 VARCHAR 列上创建倒排索引,前缀查询可以直接利用索引而不需要全表扫描:
CREATE INDEX idx_display_name ON ontology_objects(display_name)
USING INVERTED PROPERTIES("parser" = "standard");
#3.3 FUZZY — 模糊匹配
模糊搜索容忍拼写错误。用户输入 "custmer"(少了一个 o)仍然能匹配到 "Customer"。
class FuzzyMatcher:
"""基于编辑距离的模糊匹配"""
def __init__(self, max_edit_distance: int = 2):
self.max_edit_distance = max_edit_distance
def match(self, query: str, candidates: list[str]) -> list[tuple[str, int]]:
results = []
for candidate in candidates:
distance = self._levenshtein(query.lower(), candidate.lower())
if distance <= self.max_edit_distance:
results.append((candidate, distance))
return sorted(results, key=lambda x: x[1])
def _levenshtein(self, s1: str, s2: str) -> int:
if len(s1) < len(s2):
return self._levenshtein(s2, s1)
if len(s2) == 0:
return len(s1)
prev_row = range(len(s2) + 1)
for i, c1 in enumerate(s1):
curr_row = [i + 1]
for j, c2 in enumerate(s2):
insertions = prev_row[j + 1] + 1
deletions = curr_row[j] + 1
substitutions = prev_row[j] + (c1 != c2)
curr_row.append(min(insertions, deletions, substitutions))
prev_row = curr_row
return prev_row[-1]
Doris 的倒排索引也支持模糊查询——通过 MATCH_PHRASE 加 slop 参数或 MATCH_ALL 来实现近似匹配。
#3.4 EXACT — 精确匹配
精确匹配用于已知确切值的场景,比如按 ID、编码、SKU 查找:
SELECT * FROM ontology_objects
WHERE object_id = 'ORD-2024-001234'
OR attributes->>'sku' = 'ORD-2024-001234';
精确匹配绕过评分逻辑,直接返回完全匹配的结果。它是最快的模式,直接走主键或唯一索引。
#3.5 WILDCARD — 通配符匹配
通配符搜索支持 *(任意字符序列)和 ?(单个字符):
class WildcardMatcher:
"""通配符搜索转换器"""
def to_sql(self, pattern: str) -> str:
"""将通配符模式转换为 SQL LIKE 模式"""
sql_pattern = pattern.replace("*", "%").replace("?", "_")
return f"display_name LIKE '{sql_pattern}'"
def to_regex(self, pattern: str) -> str:
"""将通配符模式转换为正则表达式"""
regex = pattern.replace(".", r"\.") \
.replace("*", ".*") \
.replace("?", ".")
return f"^{regex}$"
#3.6 REGEX — 正则表达式匹配
正则搜索是最灵活但也最危险的模式。需要严格的输入校验和超时控制:
class RegexSearchHandler:
"""正则表达式搜索处理器"""
MAX_REGEX_LENGTH = 200
TIMEOUT_SECONDS = 5
DANGEROUS_PATTERNS = [
r"(.+)+", # 指数回溯
r"(a*)*", # 嵌套量词
r"(a|a)*", # 重叠交替
]
def validate(self, pattern: str) -> None:
if len(pattern) > self.MAX_REGEX_LENGTH:
raise SearchError("Regex pattern too long")
for dangerous in self.DANGEROUS_PATTERNS:
if dangerous in pattern:
raise SearchError("Potentially catastrophic regex pattern")
try:
re.compile(pattern)
except re.error as e:
raise SearchError(f"Invalid regex: {e}")
def search(self, pattern: str, field: str) -> str:
self.validate(pattern)
return f"{field} REGEXP '{pattern}'"
#3.7 模式路由器
class SearchModeRouter:
"""根据搜索模式路由到对应的处理器"""
def __init__(self):
self._handlers: dict[SearchMode, SearchHandler] = {
SearchMode.BEST_MATCH: BestMatchHandler(),
SearchMode.PREFIX: PrefixHandler(),
SearchMode.FUZZY: FuzzyHandler(),
SearchMode.EXACT: ExactHandler(),
SearchMode.WILDCARD: WildcardHandler(),
SearchMode.REGEX: RegexHandler(),
}
def route(self, request: SearchRequest) -> SearchResponse:
handler = self._handlers.get(request.mode)
if handler is None:
raise SearchError(f"Unsupported search mode: {request.mode}")
# 自动检测模式(当 mode=BEST_MATCH 时)
if request.mode == SearchMode.BEST_MATCH:
detected = self._auto_detect(request.query)
if detected != SearchMode.BEST_MATCH:
handler = self._handlers[detected]
return handler.execute(request)
def _auto_detect(self, query: str) -> SearchMode:
if query.startswith('"') and query.endswith('"'):
return SearchMode.EXACT
if "*" in query or "?" in query:
return SearchMode.WILDCARD
if query.startswith("/") and query.endswith("/"):
return SearchMode.REGEX
return SearchMode.BEST_MATCH
#4. 分面搜索(Faceted Search)
#4.1 分面类型
分面搜索让用户在搜索结果上按维度聚合,快速缩小范围。SearchService 支持三种分面类型:
| 分面类型 | 用途 | 示例 |
|---|---|---|
| TERMS | 离散值统计 | 按对象类型分组:Customer(120), Order(89), Product(45) |
| RANGE | 数值区间 | 按金额区间:0-1K(30), 1K-10K(55), 10K+(15) |
| DATE_RANGE | 时间区间 | 按创建时间:最近7天(20), 最近30天(45), 更早(35) |
#4.2 分面请求与响应
message FacetRequest {
string field = 1;
FacetType type = 2;
int32 size = 3; // TERMS: 返回前 N 个值
repeated double ranges = 4; // RANGE: 区间边界
repeated DateRange date_ranges = 5;
}
enum FacetType {
TERMS = 0;
RANGE = 1;
DATE_RANGE = 2;
}
message FacetResult {
string field = 1;
FacetType type = 2;
repeated FacetBucket buckets = 3;
}
message FacetBucket {
string key = 1;
int64 count = 2;
double from = 3; // RANGE
double to = 4; // RANGE
}
#4.3 分面查询生成
分面查询需要在不改变主查询结果的前提下,额外执行聚合统计:
class FacetQueryBuilder:
"""分面查询构建器"""
def build_terms_facet(self, field: str, size: int) -> str:
return f"""
SELECT {field} AS facet_key, COUNT(*) AS facet_count
FROM ontology_objects
WHERE {{base_where_clause}}
GROUP BY {field}
ORDER BY facet_count DESC
LIMIT {size}
"""
def build_range_facet(self, field: str, ranges: list[float]) -> str:
cases = []
for i in range(len(ranges) - 1):
lo, hi = ranges[i], ranges[i + 1]
cases.append(
f"WHEN {field} >= {lo} AND {field} < {hi} "
f"THEN '{lo}-{hi}'"
)
cases.append(f"WHEN {field} >= {ranges[-1]} THEN '{ranges[-1]}+'")
case_sql = "\n ".join(cases)
return f"""
SELECT
CASE
{case_sql}
END AS facet_key,
COUNT(*) AS facet_count
FROM ontology_objects
WHERE {{base_where_clause}}
GROUP BY facet_key
ORDER BY facet_key
"""
def build_date_range_facet(
self, field: str, ranges: list[dict]
) -> str:
cases = []
for r in ranges:
label = r["label"]
from_date = r.get("from", "1970-01-01")
to_date = r.get("to", "2099-12-31")
cases.append(
f"WHEN {field} >= '{from_date}' "
f"AND {field} < '{to_date}' "
f"THEN '{label}'"
)
case_sql = "\n ".join(cases)
return f"""
SELECT
CASE
{case_sql}
END AS facet_key,
COUNT(*) AS facet_count
FROM ontology_objects
WHERE {{base_where_clause}}
GROUP BY facet_key
"""
#4.4 分面与过滤的交互
当用户点击某个分面值进行过滤时,其他分面的统计需要更新,但被点击的分面自身不应重新过滤(否则只会显示一个选中值):
class FacetFilterCoordinator:
"""协调分面过滤与分面统计"""
def execute_with_facets(
self,
base_query: str,
facet_requests: list[FacetRequest],
active_filters: dict[str, list[str]],
) -> tuple[list[SearchHit], list[FacetResult]]:
# 主查询应用所有过滤
main_results = self._execute_main(base_query, active_filters)
# 每个分面查询排除自身的过滤
facet_results = []
for facet in facet_requests:
other_filters = {
k: v for k, v in active_filters.items()
if k != facet.field
}
facet_result = self._execute_facet(
base_query, facet, other_filters
)
facet_results.append(facet_result)
return main_results, facet_results
#5. 自动建议(Auto-Suggest)
#5.1 三层建议来源
搜索建议来自三个层次,按优先级排列:
+-------------------------------------------+
| Layer 1: Hot Queries (Redis Sorted Set) | ← 全局热门
+-------------------------------------------+
| Layer 2: Recent History (per user) | ← 个人历史
+-------------------------------------------+
| Layer 3: Entity Name Index (Trie) | ← 实体名称
+-------------------------------------------+
#5.2 Redis 热词管理
class HotQueryManager:
"""基于 Redis Sorted Set 的热词管理"""
HOT_QUERY_KEY = "search:hot_queries"
RECENT_KEY_PREFIX = "search:recent:{user_id}"
MAX_HOT_QUERIES = 1000
MAX_RECENT = 50
HOT_QUERY_TTL = 86400 * 7 # 7 天
def __init__(self, redis: Redis):
self._redis = redis
async def record_query(
self, query: str, user_id: str
) -> None:
"""记录一次搜索查询"""
normalized = query.strip().lower()
if len(normalized) < 2:
return
pipe = self._redis.pipeline()
# 热词计数 +1
pipe.zincrby(self.HOT_QUERY_KEY, 1, normalized)
# 个人历史(最近 50 条)
pipe.lpush(f"{self.RECENT_KEY_PREFIX}:{user_id}", normalized)
pipe.ltrim(f"{self.RECENT_KEY_PREFIX}:{user_id}", 0, self.MAX_RECENT - 1)
pipe.expire(f"{self.RECENT_KEY_PREFIX}:{user_id}", self.HOT_QUERY_TTL)
await pipe.execute()
# 定期清理低频词
count = await self._redis.zcard(self.HOT_QUERY_KEY)
if count > self.MAX_HOT_QUERIES:
await self._redis.zremrangebyrank(
self.HOT_QUERY_KEY, 0,
count - self.MAX_HOT_QUERIES - 1
)
async def get_suggestions(
self, prefix: str, user_id: str, limit: int = 10
) -> list[Suggestion]:
"""获取搜索建议"""
suggestions = []
# 1. 个人历史(优先)
recent = await self._redis.lrange(
f"{self.RECENT_KEY_PREFIX}:{user_id}", 0, -1
)
for q in recent:
q_str = q.decode() if isinstance(q, bytes) else q
if q_str.startswith(prefix.lower()):
suggestions.append(
Suggestion(text=q_str, source="RECENT", score=1.0)
)
if len(suggestions) >= limit // 3:
break
# 2. 全局热词
hot = await self._redis.zrevrangebyscore(
self.HOT_QUERY_KEY, "+inf", "-inf",
start=0, num=100, withscores=True
)
for q, score in hot:
q_str = q.decode() if isinstance(q, bytes) else q
if q_str.startswith(prefix.lower()):
suggestions.append(
Suggestion(text=q_str, source="HOT", score=score)
)
if len(suggestions) >= limit:
break
return suggestions[:limit]
#5.3 实体名称前缀索引
除了热词和历史,还需要基于实际数据的建议。这通过内存 Trie 树或 Doris 前缀查询实现:
class EntitySuggester:
"""基于实体名称的搜索建议"""
SUGGEST_SQL = """
SELECT display_name, object_type, object_id
FROM ontology_objects
WHERE display_name LIKE '{prefix}%'
ORDER BY view_count DESC
LIMIT {limit}
"""
async def suggest(
self, prefix: str, limit: int = 5
) -> list[Suggestion]:
rows = await self._doris.execute(
self.SUGGEST_SQL.format(prefix=prefix, limit=limit)
)
return [
Suggestion(
text=row["display_name"],
source="ENTITY",
score=0.5,
metadata={
"object_type": row["object_type"],
"object_id": row["object_id"],
},
)
for row in rows
]
#6. 保存搜索(Saved Searches)
#6.1 保存搜索模型
class SavedSearch(BaseModel):
"""保存的搜索条件"""
id: str = Field(default_factory=lambda: str(uuid4()))
user_id: str
name: str
description: str | None = None
query: str
mode: SearchMode = SearchMode.BEST_MATCH
object_types: list[str] = []
filters: list[FilterClause] = []
facet_selections: dict[str, list[str]] = {}
sort: SortSpec | None = None
is_shared: bool = False
created_at: datetime = Field(default_factory=datetime.utcnow)
updated_at: datetime = Field(default_factory=datetime.utcnow)
last_executed_at: datetime | None = None
execution_count: int = 0
#6.2 SavedSearchService
class SavedSearchService:
"""保存搜索管理服务"""
async def create(
self, user_id: str, request: CreateSavedSearchRequest
) -> SavedSearch:
saved = SavedSearch(
user_id=user_id,
name=request.name,
description=request.description,
query=request.query,
mode=request.mode,
object_types=request.object_types,
filters=request.filters,
)
await self._repo.save(saved)
return saved
async def execute(
self, saved_id: str, user_id: str
) -> SearchResponse:
saved = await self._repo.get(saved_id)
if saved.user_id != user_id and not saved.is_shared:
raise PermissionError("Cannot execute others' private search")
# 构建搜索请求
request = SearchRequest(
query=saved.query,
mode=saved.mode,
object_types=saved.object_types,
filters=saved.filters,
)
# 更新执行统计
saved.last_executed_at = datetime.utcnow()
saved.execution_count += 1
await self._repo.save(saved)
return await self._search_service.search(request)
async def list_by_user(
self, user_id: str
) -> list[SavedSearch]:
own = await self._repo.find_by_user(user_id)
shared = await self._repo.find_shared()
return own + [s for s in shared if s.user_id != user_id]
#7. 权限感知搜索过滤
#7.1 搜索权限模型
搜索结果必须遵守数据权限策略。这意味着即使 Doris 返回了匹配结果,如果用户没有该对象的读权限,结果也不能展示:
+--------------------------------------------------+
| SearchRequest |
| | |
| v |
| Query Execution (Doris) |
| | |
| v |
| Raw Results (N hits) |
| | |
| v |
| Permission Filter Layer |
| ┌──────────────────────────────────────┐ |
| │ 1. Object-type-level ACL │ |
| │ 2. Row-level security (RLS) │ |
| │ 3. Column-level masking │ |
| │ 4. World-level isolation │ |
| └──────────────────────────────────────┘ |
| | |
| Filtered Results (M hits, M <= N) |
+--------------------------------------------------+
#7.2 查询时注入权限条件
更高效的方式是在查询生成阶段就注入权限条件,避免查出大量不可见结果再过滤:
class PermissionAwareSearchFilter:
"""权限感知的搜索过滤器"""
async def inject_permission_clause(
self, user: User, base_query: str
) -> str:
# 获取用户可访问的对象类型
accessible_types = await self._acl_service.get_accessible_types(
user.id, Permission.READ
)
type_clause = (
f"object_type IN ({','.join(repr(t) for t in accessible_types)})"
)
# 获取行级安全策略
rls_policies = await self._policy_service.get_rls_policies(user.id)
rls_clauses = []
for policy in rls_policies:
rls_clauses.append(policy.to_sql_clause())
# 获取 World 隔离条件
world_clause = f"world_id = '{user.active_world_id}'"
# 组合所有权限条件
permission_clause = " AND ".join(
[type_clause, world_clause] + rls_clauses
)
return f"{base_query} AND ({permission_clause})"
async def mask_columns(
self, user: User, hits: list[SearchHit]
) -> list[SearchHit]:
"""对搜索结果应用列级脱敏"""
masking_rules = await self._policy_service.get_masking_rules(user.id)
for hit in hits:
for field, rule in masking_rules.items():
if field in hit.attributes:
hit.attributes[field] = rule.apply(
hit.attributes[field]
)
# 高亮也需要脱敏
hit.highlights = [
h for h in hit.highlights
if h.field not in masking_rules
]
return hits
#8. Doris 倒排索引集成
#8.1 为什么选择 Doris 而不是 Elasticsearch
| 对比维度 | Doris | Elasticsearch |
|---|---|---|
| 架构复杂度 | 单一系统,数据不需要同步 | 需要维护数据同步管道 |
| 数据一致性 | 强一致(同一份数据) | 最终一致(需要同步延迟) |
| 运维成本 | 低(已有 Doris 集群) | 高(额外集群 + 同步) |
| 全文检索 | 2.0+ 支持倒排索引 | 原生支持,功能更强 |
| 聚合分析 | 强项 | 也很强 |
| 中文支持 | 支持中文分词 | 需要插件 |
在智策平台的场景下,Doris 的倒排索引完全满足需求,且避免了引入 ES 带来的数据同步复杂度。
#8.2 倒排索引配置
-- 为对象表创建倒排索引
ALTER TABLE ontology_objects ADD INDEX idx_inv_display_name(display_name)
USING INVERTED
PROPERTIES(
"parser" = "chinese",
"lower_case" = "true"
);
ALTER TABLE ontology_objects ADD INDEX idx_inv_description(description)
USING INVERTED
PROPERTIES(
"parser" = "chinese",
"lower_case" = "true",
"support_phrase" = "true"
);
ALTER TABLE ontology_objects ADD INDEX idx_inv_tags(tags)
USING INVERTED
PROPERTIES(
"parser" = "comma"
);
-- 对 JSON 属性列创建倒排索引
ALTER TABLE ontology_objects ADD INDEX idx_inv_attrs(attributes)
USING INVERTED
PROPERTIES(
"parser" = "chinese",
"lower_case" = "true"
);
#8.3 全文检索查询
Doris 倒排索引支持的查询语法:
-- MATCH_ANY: 任意分词命中
SELECT * FROM ontology_objects
WHERE display_name MATCH_ANY '客户 订单';
-- MATCH_ALL: 所有分词都命中
SELECT * FROM ontology_objects
WHERE description MATCH_ALL '数据 管道 调度';
-- MATCH_PHRASE: 短语匹配(保持分词顺序)
SELECT * FROM ontology_objects
WHERE description MATCH_PHRASE '实时数据管道';
#8.4 查询构建器
class DorisSearchQueryBuilder:
"""Doris 全文检索查询构建器"""
def build(self, request: SearchRequest) -> str:
select_clause = self._build_select(request)
where_clause = self._build_where(request)
order_clause = self._build_order(request)
limit_clause = self._build_limit(request.pagination)
return f"""
{select_clause}
FROM ontology_objects
WHERE {where_clause}
{order_clause}
{limit_clause}
"""
def _build_where(self, request: SearchRequest) -> str:
clauses = []
# 搜索条件
if request.mode == SearchMode.BEST_MATCH:
fields = request.search_fields or [
"display_name", "description", "tags"
]
match_clauses = [
f"{f} MATCH_ANY '{request.query}'" for f in fields
]
clauses.append(f"({' OR '.join(match_clauses)})")
elif request.mode == SearchMode.PREFIX:
fields = request.search_fields or ["display_name"]
like_clauses = [
f"{f} LIKE '{request.query}%'" for f in fields
]
clauses.append(f"({' OR '.join(like_clauses)})")
elif request.mode == SearchMode.EXACT:
fields = request.search_fields or ["display_name"]
exact_clauses = [
f"{f} = '{request.query}'" for f in fields
]
clauses.append(f"({' OR '.join(exact_clauses)})")
# 类型过滤
if request.object_types:
types = ",".join(f"'{t}'" for t in request.object_types)
clauses.append(f"object_type IN ({types})")
# 额外过滤条件
for f in request.filters:
clauses.append(self._filter_to_sql(f))
return " AND ".join(clauses) if clauses else "1=1"
#9. 搜索高亮
#9.1 高亮实现
搜索结果中的关键词高亮帮助用户快速识别匹配位置:
class HighlightGenerator:
"""搜索结果高亮生成器"""
PRE_TAG = "<em>"
POST_TAG = "</em>"
FRAGMENT_SIZE = 150
def highlight(
self, query: str, text: str, field: str
) -> list[Highlight]:
if not text:
return []
tokens = self._tokenize(query)
fragments = self._extract_fragments(text, tokens)
return [
Highlight(
field=field,
fragment=self._mark_tokens(fragment, tokens),
)
for fragment in fragments
]
def _tokenize(self, query: str) -> list[str]:
"""分词"""
return [t.strip() for t in query.split() if t.strip()]
def _extract_fragments(
self, text: str, tokens: list[str]
) -> list[str]:
"""提取包含匹配词的文本片段"""
fragments = []
text_lower = text.lower()
for token in tokens:
pos = text_lower.find(token.lower())
if pos >= 0:
start = max(0, pos - self.FRAGMENT_SIZE // 2)
end = min(len(text), pos + len(token) + self.FRAGMENT_SIZE // 2)
fragment = text[start:end]
if start > 0:
fragment = "..." + fragment
if end < len(text):
fragment = fragment + "..."
fragments.append(fragment)
return fragments[:3] # 最多 3 个片段
def _mark_tokens(self, text: str, tokens: list[str]) -> str:
"""在文本中标记匹配词"""
result = text
for token in tokens:
pattern = re.compile(re.escape(token), re.IGNORECASE)
result = pattern.sub(
f"{self.PRE_TAG}\\g<0>{self.POST_TAG}", result
)
return result
#10. 性能优化策略
#10.1 搜索缓存
class SearchCache:
"""搜索结果缓存"""
CACHE_TTL = 300 # 5 分钟
CACHE_PREFIX = "search:cache:"
async def get_or_execute(
self, request: SearchRequest, executor: Callable
) -> SearchResponse:
cache_key = self._build_key(request)
cached = await self._redis.get(cache_key)
if cached:
return SearchResponse.model_validate_json(cached)
response = await executor(request)
# 只缓存耗时较长的查询
if response.metadata.query_time_ms > 100:
await self._redis.setex(
cache_key,
self.CACHE_TTL,
response.model_dump_json(),
)
return response
def _build_key(self, request: SearchRequest) -> str:
import hashlib
content = request.model_dump_json()
digest = hashlib.sha256(content.encode()).hexdigest()[:16]
return f"{self.CACHE_PREFIX}{digest}"
#10.2 性能指标
| 搜索模式 | 典型延迟 | 适用数据量 |
|---|---|---|
| EXACT | < 5ms | 任意 |
| PREFIX | < 20ms | < 1000 万 |
| BEST_MATCH | < 100ms | < 1000 万 |
| FUZZY | < 200ms | < 100 万 |
| WILDCARD | < 500ms | < 100 万 |
| REGEX | < 5000ms | < 10 万 |
#10.3 查询超时与熔断
class SearchCircuitBreaker:
"""搜索查询熔断器"""
def __init__(
self,
timeout_ms: int = 5000,
failure_threshold: int = 5,
reset_timeout: int = 30,
):
self._timeout_ms = timeout_ms
self._failure_count = 0
self._failure_threshold = failure_threshold
self._reset_timeout = reset_timeout
self._state = "CLOSED"
self._last_failure: datetime | None = None
async def execute(
self, coro: Coroutine
) -> SearchResponse:
if self._state == "OPEN":
if self._should_reset():
self._state = "HALF_OPEN"
else:
raise SearchError("Search service circuit breaker is OPEN")
try:
result = await asyncio.wait_for(
coro, timeout=self._timeout_ms / 1000
)
self._on_success()
return result
except asyncio.TimeoutError:
self._on_failure()
raise SearchError(f"Search timed out after {self._timeout_ms}ms")
except Exception:
self._on_failure()
raise
#11. gRPC 服务定义
service SearchService {
// 通用搜索
rpc Search(SearchRequest) returns (SearchResponse);
// 搜索建议
rpc Suggest(SuggestRequest) returns (SuggestResponse);
// 分面统计
rpc GetFacets(FacetOnlyRequest) returns (FacetOnlyResponse);
// 保存搜索
rpc CreateSavedSearch(CreateSavedSearchRequest)
returns (SavedSearch);
rpc ListSavedSearches(ListSavedSearchesRequest)
returns (ListSavedSearchesResponse);
rpc ExecuteSavedSearch(ExecuteSavedSearchRequest)
returns (SearchResponse);
rpc DeleteSavedSearch(DeleteSavedSearchRequest)
returns (google.protobuf.Empty);
}
#12. 端到端搜索流程
User types "客户 张"
|
v
AutoSuggest (< 50ms)
├── Recent: "客户 张三" (user's history)
├── Hot: "客户管理" (global hot)
└── Entity: "张伟" (Customer), "张明" (Supplier)
|
User selects "客户 张三" or presses Enter
|
v
SearchRequest(query="客户 张三", mode=BEST_MATCH)
|
v
ModeRouter → auto-detect → BEST_MATCH
|
v
PermissionFilter.inject(user) → add ACL/RLS clauses
|
v
DorisQueryBuilder.build() →
SELECT * FROM ontology_objects
WHERE (display_name MATCH_ALL '客户 张三'
OR description MATCH_ANY '客户 张三')
AND object_type IN ('Customer','Order')
AND world_id = 'main'
ORDER BY score DESC
LIMIT 20
|
v
Doris Inverted Index → raw hits
|
v
BestMatchScorer → re-rank with field boost + recency
|
v
HighlightGenerator → mark matched tokens
|
v
FacetQueryBuilder → parallel facet aggregations
|
v
SearchResponse {
hits: [{
object_type: "Customer",
object_id: "cust-001",
score: 0.95,
attributes: {display_name: "张三", ...},
highlights: [{field: "display_name",
fragment: "<em>张三</em>"}]
}, ...],
total_count: 42,
facets: [
{field: "object_type", buckets: [
{key: "Customer", count: 30},
{key: "Order", count: 12}
]}
]
}
|
v
HotQueryManager.record("客户 张三", user_id)
#Key Takeaways
- 6 种搜索模式覆盖所有场景 — 从精确查找到正则表达式,每种模式有独立的优化路径
- 分面搜索不是后处理 — 分面统计在查询时并行执行,且遵循"排除自身过滤"原则
- 三层建议系统 — 个人历史 > 全局热词 > 实体名称,Redis Sorted Set 实现 O(log N) 更新
- 权限注入在查询时 — 不是查出结果再过滤,而是将权限条件注入 SQL,减少无效查询
- Doris 倒排索引替代 ES — 单一系统解决 OLAP + 全文检索,避免数据同步的一致性问题
- 熔断和缓存保障稳定性 — 正则搜索等高风险模式有超时和熔断保护
#Next Article
下一篇 S3-14 指标系统:6 种计算策略的优先级路由 将深入指标注册、计算策略路由和 OQL 重写器的实现。
Tags: #SearchEngine #FacetedSearch #InvertedIndex #Doris #AutoSuggest #Redis #PermissionAware #gRPC #OntologyPlatform