返回博客

源码精读:AnalyticsQueryService — 14 种聚合的实现

DefaultAnalyticsQueryService 是 Data Layer 中专注于聚合分析的服务,基于 Quarkus 3.x 实现。它通过四个独立的 SQL Builder(AggregateGroupedSqlBuilder、TimeBucketSqlBuilder、TopNSqlBuilder、DistributionSqlBuilder)将 Proto 请求转换为参数化 Doris SQL,支持 14 种聚合函数(SUM/COUNT/AVG/MIN/MAX/PERCENTILE 等)、ROLLUP 小计行检测、Redis 30 秒结果缓存、时间桶空桶填充(BucketFiller),以及分布查询的等宽/等频/自定义三种分桶模式。本文将逐行分析其分层架构、SQL Builder 的安全参数化策略、缓存键的 MD5 生成、ROLLUP 小计行的 NULL 检测、时间桶的空桶填充算法,以及分布查询的三种分桶实现。

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

源码精读:AnalyticsQueryService — 14 种聚合的实现

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

#TL;DR

DefaultAnalyticsQueryService 是 Data Layer 中专注于聚合分析的服务,基于 Quarkus 3.x 实现。它通过四个独立的 SQL Builder(AggregateGroupedSqlBuilderTimeBucketSqlBuilderTopNSqlBuilderDistributionSqlBuilder)将 Proto 请求转换为参数化 Doris SQL,支持 14 种聚合函数(SUM/COUNT/AVG/MIN/MAX/PERCENTILE 等)、ROLLUP 小计行检测、Redis 30 秒结果缓存、时间桶空桶填充(BucketFiller),以及分布查询的等宽/等频/自定义三种分桶模式。本文将逐行分析其分层架构、SQL Builder 的安全参数化策略、缓存键的 MD5 生成、ROLLUP 小计行的 NULL 检测、时间桶的空桶填充算法,以及分布查询的三种分桶实现。

#目录

  1. 整体架构与四种查询类型
  2. 分组聚合:aggregateGrouped 的 8 步流程
  3. SQL Builder 模式:AggregateGroupedSqlBuilder
  4. Redis 缓存策略:MD5 键与 30 秒 TTL
  5. ROLLUP 小计行检测
  6. 时间桶聚合:TimeBucketSqlBuilder 与空桶填充
  7. TopN 查询:排序与分页
  8. 分布查询:三种分桶模式
  9. QueryExecutionInfo:执行信息透传
  10. 14 种聚合函数清单
  11. Key Takeaways

#1. 整体架构与四种查询类型

Java
@ApplicationScoped
public class DefaultAnalyticsQueryService implements AnalyticsQueryService {
    private final DorisClient dorisClient;
    private final RedisDataSource redisDataSource;
    private final ObjectMapper objectMapper;
}

四种查询类型对应四个接口方法:

方法用途SQL Builder典型场景
aggregateGrouped分组聚合AggregateGroupedSqlBuilder按部门统计销售额
aggregateTimeBucketed时间桶聚合TimeBucketSqlBuilder按小时/天/月趋势
queryTopNTopN 排名TopNSqlBuilder销售额前 10 名
queryDistribution分布统计DistributionSqlBuilder年龄分布直方图

#2. 分组聚合:aggregateGrouped 的 8 步流程

Java
@Override
public AggregateGroupedResponse aggregateGrouped(AggregateGroupedRequest request) {
    long startMs = System.currentTimeMillis();
    String worldId = extractWorldId(request);

    // 1. 缓存查找
    if (request.getUseCache()) {
        String cacheKey = buildCacheKey(worldId, objectType, request.toString());
        AggregateGroupedResponse cached = readFromCache(cacheKey);
        if (cached != null) return cached;
    }

    // 2. 构建 SQL
    AggregateGroupedSqlBuilder builder = new AggregateGroupedSqlBuilder(request);
    AggregateGroupedSqlBuilder.BuildResult buildResult = builder.build();

    // 3. 执行主查询
    List<Map<String, Object>> rows = dorisClient.executeQuery(
            buildResult.getSql(), buildResult.getParams().toArray());

    // 4. 映射行
    List<AggregateRow> aggregateRows = mapRows(rows, request);

    // 5. 总计行(独立查询,无 GROUP BY)
    AggregateRow totalRow = null;
    if (request.getIncludeTotal() && buildResult.hasTotalQuery()) {
        List<Map<String, Object>> totalRows = dorisClient.executeQuery(
                totalQuery.getSql(), totalQuery.getParams().toArray());
        if (!totalRows.isEmpty()) {
            totalRow = mapSingleRow(totalRows.get(0), request, false);
        }
    }

    // 6. 执行信息
    QueryExecutionInfo executionInfo = QueryExecutionInfo.newBuilder()
            .setExecuteMs(executeMs).setTotalMs(totalMs)
            .setCacheHit(false).setQueryId(UUID.randomUUID().toString()).build();

    // 7. 组装响应
    AggregateGroupedResponse response = responseBuilder
            .addAllRows(aggregateRows).setExecutionInfo(executionInfo).build();

    // 8. 写入缓存
    if (request.getUseCache()) writeToCache(cacheKey, response);

    return response;
}

设计亮点

  • 缓存前置:缓存查找在 SQL 构建之前执行,命中时零 SQL 开销
  • 总计行独立查询:总计行使用单独的 SQL(无 GROUP BY),确保计算独立性
  • 执行信息透传QueryExecutionInfo 包含 executeMstotalMscacheHitqueryId

#3. SQL Builder 模式:AggregateGroupedSqlBuilder

SQL Builder 将 Proto 请求转换为参数化 SQL,核心职责:

  1. AggregateColumn 列表提取聚合函数和字段名
  2. groupByFields 构建 GROUP BY 子句
  3. filters 构建参数化 WHERE 子句
  4. 可选添加 WITH ROLLUP 生成小计行

所有用户输入通过 ? 占位符传递,防止 SQL 注入。

Java
AggregateGroupedSqlBuilder builder = new AggregateGroupedSqlBuilder(request);
AggregateGroupedSqlBuilder.BuildResult buildResult = builder.build();

// buildResult.getSql() => "SELECT dept, SUM(revenue) FROM ... WHERE world_id = ? GROUP BY dept"
// buildResult.getParams() => ["world_main"]

#4. Redis 缓存策略:MD5 键与 30 秒 TTL

Java
private static final int CACHE_TTL_SECONDS = 30;
private static final String CACHE_KEY_PREFIX = "analytics:grouped:";

private String buildCacheKey(String worldId, String objectType, String requestStr) {
    String raw = worldId + objectType + requestStr;
    return CACHE_KEY_PREFIX + md5(raw);
}

30 秒 TTL 的理由:分析查询的结果在短时间内通常不会变化,30 秒缓存在"实时性"和"重复查询性能"之间取得平衡。

MD5 键:使用 worldId + objectType + request.toString() 的 MD5 摘要作为缓存键,避免超长键名。

#5. ROLLUP 小计行检测

WITH ROLLUP 启用时,Doris 返回额外的小计行。小计行的特征是 GROUP BY 字段值为 NULL

Code
department | region | SUM(revenue)
-----------|--------|-------------
Sales      | East   | 1000         <- 普通行
Sales      | West   | 2000         <- 普通行
Sales      | NULL   | 3000         <- 小计行(Sales 合计)
NULL       | NULL   | 5000         <- 总计行

当至少一个 GROUP BY 字段为 null 时,该行被标记为 is_subtotal = true

#6. 时间桶聚合:TimeBucketSqlBuilder 与空桶填充

时间桶聚合将连续时间轴划分为固定大小的桶(小时/天/周/月),对每个桶执行聚合。

Java
// TimeBucketSqlBuilder 生成:
// SELECT DATE_FORMAT(created_at, '%Y-%m-%d') AS bucket, SUM(amount)
// FROM orders WHERE world_id = ? AND created_at BETWEEN ? AND ?
// GROUP BY bucket ORDER BY bucket

BucketFiller 负责填充没有数据的空桶:

填充策略行为适用场景
ZERO空桶填 0计数、求和类指标
PREVIOUS继承前一个桶的值累计类指标
NONE不填充,空桶缺失稀疏数据展示

#7. TopN 查询:排序与分页

TopN 查询返回按指定度量字段排序的前 N 条记录:

Java
// TopNSqlBuilder 生成:
// SELECT entity_id, name, SUM(revenue) AS measure
// FROM ... WHERE world_id = ?
// GROUP BY entity_id, name
// ORDER BY measure DESC LIMIT ?

支持 ASC/DESC 排序方向和可选的过滤条件。

#8. 分布查询:三种分桶模式

模式说明SQL 策略
EQUAL_WIDTH等宽分桶FLOOR(value / width)
EQUAL_FREQUENCY等频分桶NTILE() 窗口函数
CUSTOM_BOUNDARIES自定义边界CASE WHEN 表达式

分布查询还会计算统计信息(DistributionStats):min、max、mean、median、stddev,为直方图提供辅助上下文。

#9. QueryExecutionInfo:执行信息透传

Java
QueryExecutionInfo executionInfo = QueryExecutionInfo.newBuilder()
        .setExecuteMs(executeMs)     // SQL 执行耗时
        .setTotalMs(totalMs)         // 总处理耗时(含缓存查找)
        .setCacheHit(false)          // 是否缓存命中
        .setQueryId(UUID.randomUUID().toString())  // 查询追踪 ID
        .build();

每个分析查询响应都包含执行信息,支持前端展示"查询耗时"和"是否缓存命中"。

#10. 14 种聚合函数清单

#函数说明Doris SQL
1COUNT计数COUNT(*) / COUNT(field)
2COUNT_DISTINCT去重计数COUNT(DISTINCT field)
3SUM求和SUM(field)
4AVG平均AVG(field)
5MIN最小值MIN(field)
6MAX最大值MAX(field)
7STDDEV标准差STDDEV(field)
8VARIANCE方差VARIANCE(field)
9PERCENTILE_50中位数PERCENTILE_APPROX(field, 0.5)
10PERCENTILE_90P90PERCENTILE_APPROX(field, 0.9)
11PERCENTILE_95P95PERCENTILE_APPROX(field, 0.95)
12PERCENTILE_99P99PERCENTILE_APPROX(field, 0.99)
13FIRST_VALUE首值FIRST_VALUE(field)
14LAST_VALUE末值LAST_VALUE(field)

#11. Key Takeaways

  1. SQL Builder 模式:每种查询类型有独立的 SQL Builder,将 Proto 请求安全转换为参数化 SQL
  2. Redis 30 秒缓存:MD5 缓存键 + 30 秒 TTL,在实时性和性能之间平衡
  3. ROLLUP 小计检测:通过 GROUP BY 字段的 NULL 值识别小计行
  4. 空桶填充BucketFiller 支持 ZERO/PREVIOUS/NONE 三种填充策略
  5. 14 种聚合:从基础的 COUNT/SUM 到高级的 PERCENTILE_APPROX/STDDEV
  6. 执行信息透传QueryExecutionInfo 携带耗时和缓存状态,支持前端监控

#下一篇

S9-08:SearchService — 6 种搜索模式的统一抽象。我们将深入全文搜索服务,看它如何在 Doris OLAP 引擎上实现多种搜索模式以及 Redis 驱动的最近/热门搜索。

Tags: #coomia-dip #source-code-reading #data-Layer #analytics #aggregation #doris #sql-builder #redis-cache #rollup #distribution