返回博客

源码精读:QueryFederationService — 多引擎查询路由

QueryFederationGrpcService 是 Data Layer(Data Layer)中查询联邦层的核心服务,基于 Quarkus 3.x + gRPC 实现。它接收 OQL 查询,经过解析 → 优化 → 路由 → 执行的四阶段流水线,将查询路由到 Doris(OLAP)或 DuckDB(嵌入式分析)执行,支持缓存加速、流式返回、查询计划解释、向量搜索和图遍历五种查询模式。本文将深入分析其八依赖注入架构、查询执行的完整流水线、StorageRouter 的联邦判定逻辑、缓存策略的聚合感知、VirtualPropertyEnricher 的派生属性填充,以及异步查询状态管理。

Coomia发布于 2025年12月4日10 分钟阅读
分享本文Twitter / X

源码精读:QueryFederationService — 多引擎查询路由

系列:S9 源码精读 · 第 5 篇 | 难度:高级 | 阅读时间:25 分钟

#TL;DR

QueryFederationGrpcService 是 Data Layer(Data Layer)中查询联邦层的核心服务,基于 Quarkus 3.x + gRPC 实现。它接收 OQL 查询,经过解析 → 优化 → 路由 → 执行的四阶段流水线,将查询路由到 Doris(OLAP)或 DuckDB(嵌入式分析)执行,支持缓存加速、流式返回、查询计划解释、向量搜索和图遍历五种查询模式。本文将深入分析其八依赖注入架构、查询执行的完整流水线、StorageRouter 的联邦判定逻辑、缓存策略的聚合感知、VirtualPropertyEnricher 的派生属性填充,以及异步查询状态管理。

#目录

  1. 整体架构与八大协作者
  2. 查询执行流水线:executeQuery 的 8 步骤
  3. StorageRouter:联邦查询判定
  4. 缓存策略:聚合感知的查询缓存
  5. VirtualPropertyEnricher:派生属性填充
  6. 流式查询:executeQueryStream
  7. 查询验证与计划解释:validateQuery / explainQuery
  8. 向量搜索与图遍历:OQL 统一抽象
  9. 时间旅行查询:预留接口设计
  10. 异步查询状态管理
  11. Key Takeaways

#1. 整体架构与八大协作者

Java
@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 步骤

Java
@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 状态码含义
OQLSyntaxExceptionINVALID_ARGUMENTOQL 语法错误(客户端问题)
OQLSemanticExceptionINVALID_ARGUMENTOQL 语义错误(类型不匹配等)
QueryExecutionExceptionINTERNAL查询执行失败(服务端问题)

#3. StorageRouter:联邦查询判定

Java
if (router.requiresFederation(plan)) {
    result = router.executeFederatedQuery(plan, context);
} else {
    result = executor.execute(plan, context);
}

StorageRouter 分析物理执行计划,判断查询是否需要跨引擎联邦:

  • Doris:大规模 OLAP 查询、聚合、全文搜索
  • DuckDB:小数据量嵌入式分析、临时查询
  • 联邦:跨两个引擎的 JOIN 或 UNION 查询
Java
// 查询计划解释
var routingDecision = router.explain(plan, context);
for (var storage : routingDecision.involvedStorages()) {
    planBuilder.addEngines(storage.code());
}

explainQuery 方法输出的 engines 字段告诉调用者查询将涉及哪些存储引擎——这对性能预估和调试非常有价值。

#4. 缓存策略:聚合感知的查询缓存

Java
boolean isAggregation = plan instanceof PhysicalAggregatePlan;
boolean shouldCache = options.enableCache()
        && cacheService.isEnabled()
        && (!isAggregation || cacheService.isCacheAggregationsEnabled());

聚合查询的缓存特殊处理:聚合查询的结果可能随数据变化而过时,因此通过 isCacheAggregationsEnabled() 提供独立开关。默认情况下,普通查询可缓存,聚合查询的缓存需要显式启用。

缓存键生成

Java
Map<String, Object> params = extractQueryParams(request);
cacheKey = cacheService.generateCacheKey(worldId, query, params);

缓存键包含三个维度:worldId(数据隔离)、query(查询文本)、params(分页参数等)。这确保了不同世界、不同查询、不同分页的结果不会互相污染。

#5. VirtualPropertyEnricher:派生属性填充

Java
// 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

Java
@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 查询验证

Java
@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();
}

精确错误定位:语法错误返回 linecolumn,IDE/编辑器可以据此定位具体的出错位置——这是开发者体验的关键。

#7.2 查询计划解释

Java
@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 向量搜索

Java
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 图遍历

Java
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. 时间旅行查询:预留接口设计

Java
@Override
public void queryAtTimestamp(QueryAtTimestampRequest request,
        StreamObserver<ExecuteQueryResponse> responseObserver) {
    responseObserver.onError(Status.UNIMPLEMENTED
            .withDescription("Time travel queries not yet implemented")
            .asRuntimeException());
}

三个时间旅行方法(queryAtTimestampqueryAtVersionqueryDiff)目前返回 UNIMPLEMENTED。这是 gRPC 接口预留的标准做法——Proto 定义和 gRPC 桩代码已就绪,等待 Nessie 集成完成后实现。

#10. 异步查询状态管理

Java
private record AsyncQueryState(
    QueryStatus status,
    QueryResult result,
    String error
) {}

private enum QueryStatus {
    PENDING, RUNNING, COMPLETED, FAILED, CANCELLED
}

AsyncQueryStateConcurrentHashMap 管理,支持长时间运行查询的状态追踪。五个状态形成完整的生命周期:PENDING → RUNNING → COMPLETED/FAILED,以及任何阶段可触发的 CANCELLED

#11. Key Takeaways

  1. 四阶段流水线:解析(OQL → AST)→ 优化(AST → Plan)→ 路由(Doris/DuckDB/联邦)→ 执行,每个阶段由独立协作者负责
  2. 聚合感知缓存:普通查询默认可缓存,聚合查询需通过独立开关启用,避免过时数据
  3. OQL 统一抽象:向量搜索(SIMILAR_TO)和图遍历(CONNECTED_TO)都转换为 OQL 查询,复用完整流水线
  4. 派生属性填充:在执行后、缓存前插入 VirtualPropertyEnricher,零虚拟属性时零开销
  5. 流式返回executeQueryStream 使用 gRPC 服务端流逐行发送,适用于大结果集
  6. 精确错误定位:语法错误返回行号和列号,支持 IDE 级错误定位
  7. 接口预留:时间旅行查询返回 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