源码精读:QueryFederationService — 多引擎查询路由
QueryFederationGrpcService 是 Data Layer(Data Layer)中查询联邦层的核心服务,基于 Quarkus 3.x + gRPC 实现。它接收 OQL 查询,经过解析 → 优化 → 路由 → 执行的四阶段流水线,将查询路由到 Doris(OLAP)或 DuckDB(嵌入式分析)执行,支持缓存加速、流式返回、查询计划解释、向量搜索和图遍历五种查询模式。本文将深入分析其八依赖注入架构、查询执行的完整流水线、StorageRouter 的联邦判定逻辑、缓存策略的聚合感知、VirtualPropertyEnricher 的派生属性填充,以及异步查询状态管理。
源码精读:QueryFederationService — 多引擎查询路由
“系列:S9 源码精读 · 第 5 篇 | 难度:高级 | 阅读时间:25 分钟
#TL;DR
QueryFederationGrpcService 是 Data Layer(Data Layer)中查询联邦层的核心服务,基于 Quarkus 3.x + gRPC 实现。它接收 OQL 查询,经过解析 → 优化 → 路由 → 执行的四阶段流水线,将查询路由到 Doris(OLAP)或 DuckDB(嵌入式分析)执行,支持缓存加速、流式返回、查询计划解释、向量搜索和图遍历五种查询模式。本文将深入分析其八依赖注入架构、查询执行的完整流水线、StorageRouter 的联邦判定逻辑、缓存策略的聚合感知、VirtualPropertyEnricher 的派生属性填充,以及异步查询状态管理。
#目录
- 整体架构与八大协作者
- 查询执行流水线:executeQuery 的 8 步骤
- StorageRouter:联邦查询判定
- 缓存策略:聚合感知的查询缓存
- VirtualPropertyEnricher:派生属性填充
- 流式查询:executeQueryStream
- 查询验证与计划解释:validateQuery / explainQuery
- 向量搜索与图遍历:OQL 统一抽象
- 时间旅行查询:预留接口设计
- 异步查询状态管理
- Key Takeaways
#1. 整体架构与八大协作者
@GrpcService
public class QueryFederationGrpcService
extends QueryFederationServiceGrpc.QueryFederationServiceImplBase {
private final OQLParserService parser;
private final QueryOptimizer optimizer;
private final QueryExecutor executor;
private final StorageRouter router;
private final QueryProtoConverter queryConverter;
private final ResultProtoConverter resultConverter;
private final QueryCacheService cacheService;
private final VirtualPropertyEnricher virtualPropertyEnricher;
// Async query state cache
private final Map<String, AsyncQueryState> asyncQueries = new ConcurrentHashMap<>();
}
八大协作者分工:
| 协作者 | 阶段 | 职责 |
|---|---|---|
OQLParserService | 解析 | OQL 文本 → AST |
QueryOptimizer | 优化 | AST → 物理执行计划 |
QueryExecutor | 执行 | 单引擎查询执行 |
StorageRouter | 路由 | 联邦查询判定与跨引擎执行 |
QueryProtoConverter | 入口转换 | gRPC 请求 → 查询上下文 |
ResultProtoConverter | 出口转换 | 查询结果 → gRPC 响应 |
QueryCacheService | 缓存 | 查询结果的 Redis 缓存 |
VirtualPropertyEnricher | 增强 | 填充 Reasoning & Decision Layer 计算的派生属性 |
这是整个 coomia-dip 中依赖注入最多的服务之一,反映了查询联邦的复杂性。
#2. 查询执行流水线:executeQuery 的 8 步骤
@Override
public void executeQuery(ExecuteQueryRequest request,
StreamObserver<ExecuteQueryResponse> responseObserver) {
String worldId = queryConverter.extractWorldId(request.getContext());
String query = request.getQuery();
try {
long startTime = System.currentTimeMillis();
// Step 1: 构建查询上下文
var context = queryConverter.toQueryContext(request);
var options = context.getOptions();
// Step 2: 缓存检查
String cacheKey = null;
if (options.enableCache() && cacheService.isEnabled()) {
cacheKey = cacheService.generateCacheKey(worldId, query,
extractQueryParams(request));
var cachedResult = cacheService.get(cacheKey);
if (cachedResult.isPresent()) {
// 缓存命中 — 直接返回
responseObserver.onNext(resultConverter.toProto(
QueryResult.of(cachedResult.get().rows(), metadata)));
responseObserver.onCompleted();
return;
}
}
// Step 3: 解析 OQL → AST
var ast = parser.parse(query);
// Step 4: 优化 AST → 物理执行计划
var plan = optimizer.optimize(ast, context);
// Step 5: 聚合缓存策略判定
boolean isAggregation = plan instanceof PhysicalAggregatePlan;
boolean shouldCache = options.enableCache()
&& cacheService.isEnabled()
&& (!isAggregation || cacheService.isCacheAggregationsEnabled());
// Step 6: 路由与执行
QueryResult result;
if (router.requiresFederation(plan)) {
result = router.executeFederatedQuery(plan, context);
} else {
result = executor.execute(plan, context);
}
// Step 6b: 派生属性填充(FEAT-005)
String ontologyType = ast.from().entityType();
result = virtualPropertyEnricher.enrich(result, context, ontologyType);
// Step 7: 写入缓存
if (shouldCache && cacheKey != null) {
cacheService.put(cacheKey, cachedResult, options.cacheTtlSeconds());
}
// Step 8: 转换并返回响应
responseObserver.onNext(resultConverter.toProto(result));
responseObserver.onCompleted();
} catch (OQLSyntaxException e) {
responseObserver.onError(Status.INVALID_ARGUMENT
.withDescription("OQL syntax error: " + e.getMessage())...);
} catch (OQLSemanticException e) {
responseObserver.onError(Status.INVALID_ARGUMENT
.withDescription("OQL semantic error: " + e.getMessage())...);
} catch (QueryExecutionException e) {
responseObserver.onError(Status.INTERNAL
.withDescription("Query execution failed: " + e.getMessage())...);
}
}
异常分层映射:
| 异常类型 | gRPC 状态码 | 含义 |
|---|---|---|
OQLSyntaxException | INVALID_ARGUMENT | OQL 语法错误(客户端问题) |
OQLSemanticException | INVALID_ARGUMENT | OQL 语义错误(类型不匹配等) |
QueryExecutionException | INTERNAL | 查询执行失败(服务端问题) |
#3. StorageRouter:联邦查询判定
if (router.requiresFederation(plan)) {
result = router.executeFederatedQuery(plan, context);
} else {
result = executor.execute(plan, context);
}
StorageRouter 分析物理执行计划,判断查询是否需要跨引擎联邦:
- Doris:大规模 OLAP 查询、聚合、全文搜索
- DuckDB:小数据量嵌入式分析、临时查询
- 联邦:跨两个引擎的 JOIN 或 UNION 查询
// 查询计划解释
var routingDecision = router.explain(plan, context);
for (var storage : routingDecision.involvedStorages()) {
planBuilder.addEngines(storage.code());
}
explainQuery 方法输出的 engines 字段告诉调用者查询将涉及哪些存储引擎——这对性能预估和调试非常有价值。
#4. 缓存策略:聚合感知的查询缓存
boolean isAggregation = plan instanceof PhysicalAggregatePlan;
boolean shouldCache = options.enableCache()
&& cacheService.isEnabled()
&& (!isAggregation || cacheService.isCacheAggregationsEnabled());
聚合查询的缓存特殊处理:聚合查询的结果可能随数据变化而过时,因此通过 isCacheAggregationsEnabled() 提供独立开关。默认情况下,普通查询可缓存,聚合查询的缓存需要显式启用。
缓存键生成:
Map<String, Object> params = extractQueryParams(request);
cacheKey = cacheService.generateCacheKey(worldId, query, params);
缓存键包含三个维度:worldId(数据隔离)、query(查询文本)、params(分页参数等)。这确保了不同世界、不同查询、不同分页的结果不会互相污染。
#5. VirtualPropertyEnricher:派生属性填充
// Step 6b: Enrich with virtual properties (FEAT-005)
String ontologyType = ast.from().entityType();
result = virtualPropertyEnricher.enrich(result, context, ontologyType);
VirtualPropertyEnricher 在查询执行后、缓存写入前执行,填充由 Reasoning & Decision Layer DerivedPropertyEngine 计算的虚拟属性。注释中的"Zero overhead if no virtual properties are defined for this type"说明了性能考量:如果没有定义虚拟属性,这步操作的开销为零。
#6. 流式查询:executeQueryStream
@Override
public void executeQueryStream(ExecuteQueryRequest request,
StreamObserver<QueryResultRow> responseObserver) {
try {
var ast = parser.parse(request.getQuery());
var context = queryConverter.toQueryContext(request);
var plan = optimizer.optimize(ast, context);
var result = executor.execute(plan, context);
// 逐行流式发送
long rowIndex = 0;
for (var row : result.rows()) {
var protoRow = resultConverter.toProtoResultRow(row, rowIndex++);
responseObserver.onNext(protoRow);
}
responseObserver.onCompleted();
} catch (Exception e) {
responseObserver.onError(Status.INTERNAL.withDescription(e.getMessage())...);
}
}
Server Streaming RPC:与 executeQuery 返回完整结果不同,executeQueryStream 使用 gRPC 服务端流,逐行发送 QueryResultRow。适用于大结果集场景——客户端无需等待所有数据加载到内存。
#7. 查询验证与计划解释:validateQuery / explainQuery
#7.1 查询验证
@Override
public void validateQuery(ValidateQueryRequest request,
StreamObserver<ValidateQueryResponse> responseObserver) {
var responseBuilder = ValidateQueryResponse.newBuilder();
try {
var ast = parser.parse(request.getQuery());
responseBuilder.setValid(true)
.setStructure(QueryStructure.newBuilder()
.setQueryType("SELECT")
.addFromTypes(ast.from().entityType())
.build());
} catch (OQLSyntaxException e) {
responseBuilder.setValid(false)
.addErrors(QueryValidationError.newBuilder()
.setErrorType("SYNTAX_ERROR")
.setMessage(e.getMessage())
.setLine(firstError.line())
.setColumn(firstError.column())
.build());
}
responseObserver.onNext(responseBuilder.build());
responseObserver.onCompleted();
}
精确错误定位:语法错误返回 line 和 column,IDE/编辑器可以据此定位具体的出错位置——这是开发者体验的关键。
#7.2 查询计划解释
@Override
public void explainQuery(ExplainQueryRequest request,
StreamObserver<ExplainQueryResponse> responseObserver) {
var ast = parser.parse(request.getQuery());
var plan = optimizer.optimize(ast, context);
var planBuilder = QueryPlan.newBuilder()
.setPlanText(plan.explain())
.setUsesIndexes(false);
// 路由决策
var routingDecision = router.explain(plan, context);
for (var storage : routingDecision.involvedStorages()) {
planBuilder.addEngines(storage.code());
}
// 代价估算
var cost = QueryCostEstimate.newBuilder()
.setEstimatedResultRows(plan.estimatedRows())
.setCostScore(plan.estimatedCost())
.build();
responseObserver.onNext(ExplainQueryResponse.newBuilder()
.setPlan(planBuilder.build())
.setCost(cost)
.build());
}
三维度输出:explainQuery 返回执行计划文本(planText)、涉及的存储引擎(engines)和代价估算(estimatedRows + costScore),为查询优化提供全方位信息。
#8. 向量搜索与图遍历:OQL 统一抽象
#8.1 向量搜索
private String buildVectorSearchOQL(VectorSearchRequest request) {
var sb = new StringBuilder("SELECT * FROM ");
sb.append(request.getSchemaTypes(0));
sb.append(" WHERE SIMILAR_TO(")
.append(request.getVectorField())
.append(", [").append(vectorStr).append("], ")
.append(request.getMinSimilarity()).append(")");
sb.append(" LIMIT ").append(request.getTopK());
return sb.toString();
}
OQL 统一抽象:向量搜索不是独立的执行路径——它被转换为带 SIMILAR_TO 函数的 OQL 查询,复用整个解析 → 优化 → 路由 → 执行流水线。
#8.2 图遍历
private String buildGraphTraversalOQL(GraphTraversalRequest request) {
var sb = new StringBuilder("SELECT * FROM ");
sb.append(pattern.getTargetTypes(0));
sb.append(" WHERE CONNECTED_TO('")
.append(request.getStartIds(0))
.append("', '").append(pattern.getRelationTypes(0))
.append("', 1, ").append(request.getMaxDepth()).append(")");
sb.append(" LIMIT ").append(request.getMaxResults());
return sb.toString();
}
图遍历同样转换为 CONNECTED_TO 函数的 OQL 查询。CONNECTED_TO(startId, relationType, minDepth, maxDepth) 四参数模型覆盖了从直接关系到多跳遍历的所有场景。
#9. 时间旅行查询:预留接口设计
@Override
public void queryAtTimestamp(QueryAtTimestampRequest request,
StreamObserver<ExecuteQueryResponse> responseObserver) {
responseObserver.onError(Status.UNIMPLEMENTED
.withDescription("Time travel queries not yet implemented")
.asRuntimeException());
}
三个时间旅行方法(queryAtTimestamp、queryAtVersion、queryDiff)目前返回 UNIMPLEMENTED。这是 gRPC 接口预留的标准做法——Proto 定义和 gRPC 桩代码已就绪,等待 Nessie 集成完成后实现。
#10. 异步查询状态管理
private record AsyncQueryState(
QueryStatus status,
QueryResult result,
String error
) {}
private enum QueryStatus {
PENDING, RUNNING, COMPLETED, FAILED, CANCELLED
}
AsyncQueryState 用 ConcurrentHashMap 管理,支持长时间运行查询的状态追踪。五个状态形成完整的生命周期:PENDING → RUNNING → COMPLETED/FAILED,以及任何阶段可触发的 CANCELLED。
#11. Key Takeaways
- 四阶段流水线:解析(OQL → AST)→ 优化(AST → Plan)→ 路由(Doris/DuckDB/联邦)→ 执行,每个阶段由独立协作者负责
- 聚合感知缓存:普通查询默认可缓存,聚合查询需通过独立开关启用,避免过时数据
- OQL 统一抽象:向量搜索(
SIMILAR_TO)和图遍历(CONNECTED_TO)都转换为 OQL 查询,复用完整流水线 - 派生属性填充:在执行后、缓存前插入
VirtualPropertyEnricher,零虚拟属性时零开销 - 流式返回:
executeQueryStream使用 gRPC 服务端流逐行发送,适用于大结果集 - 精确错误定位:语法错误返回行号和列号,支持 IDE 级错误定位
- 接口预留:时间旅行查询返回
UNIMPLEMENTED,Proto 和桩代码已就绪
#下一篇
S9-06:OQL Parser — 从文本到执行计划。我们将深入 OQL 的手写词法分析器和递归下降解析器,看它如何将查询文本转换为 AST,支持 30+ 关键字和 6 种条件类型。
Tags: #coomia-dip #source-code-reading #data-Layer #query-federation #doris #duckdb #oql #caching #vector-search #graph-traversal