共计 2066 个字符,预计需要花费 6 分钟才能阅读完成。
背景痛点
在处理 ClickHouse Rollup 查询结果时,我们经常遇到需要将扁平化的数据聚合成父子层级结构的场景。传统做法通常有两种:

- 多次查询数据库获取层级关系
- 全量加载数据后在内存中处理
这两种方式都存在明显缺陷。多次查询会导致网络 I / O 开销成倍增加,特别是在数据量大时,响应时间会变得不可接受。而全量加载方式虽然减少了查询次数,但内存消耗大,处理效率低,特别是在构建复杂层级关系时,时间复杂度可能达到 O(n²)。
技术方案对比
针对这个问题,我们评估了几种常见的技术方案:
- 递归算法 :实现简单但性能差,栈溢出风险高
- 内存索引构建 :预处理数据关系,查询效率高
- MapReduce:适合分布式处理但引入额外复杂度
经过对比,我们选择了基于内存索引构建的方案,它能够在单机环境下提供最佳的性能表现,同时保持代码的可维护性。
核心实现
内存索引结构
我们使用多级 Map 构建内存索引,第一级 Map 以父 ID 为 key,第二级 Map 维护父子关系。这种结构可以将层级关系的查询时间复杂度降低到 O(1)。
Map<String, Map<String, Node>> index = new ConcurrentHashMap<>();
并行流处理
利用 Java 的并行流特性,我们实现了数据的并行预处理和聚合:
rollupData.parallelStream().forEach(node -> {index.computeIfAbsent(node.getParentId(), k -> new ConcurrentHashMap<>())
.put(node.getId(), node);
});
边界情况处理
对于循环引用等异常情况,我们通过添加访问标志和检测机制来避免无限递归:
if (visitedNodes.contains(nodeId)) {throw new IllegalArgumentException("检测到循环引用:" + nodeId);
}
完整代码示例
数据模型定义
@Data
public class Node {
private String id;
private String parentId;
private String name;
private BigDecimal value;
private List<Node> children = new ArrayList<>();}
索引构建逻辑
public Map<String, Node> buildHierarchy(List<Node> flatData) {
// 构建索引
Map<String, Node> index = new HashMap<>();
Map<String, List<Node>> childrenMap = new HashMap<>();
flatData.forEach(node -> {index.put(node.getId(), node);
childrenMap.computeIfAbsent(node.getParentId(), k -> new ArrayList<>())
.add(node);
});
// 构建层级
childrenMap.forEach((parentId, children) -> {Node parent = index.get(parentId);
if (parent != null) {parent.setChildren(children);
}
});
return index;
}
聚合算法实现
public Node aggregateHierarchy(Node root, Map<String, Node> index) {if (root == null) return null;
// 聚合子节点
BigDecimal sum = root.getValue();
for (Node child : root.getChildren()) {Node aggregated = aggregateHierarchy(child, index);
sum = sum.add(aggregated.getValue());
}
root.setValue(sum);
return root;
}
性能考量
我们对方案进行了全面的性能测试:
- 时间复杂度分析 :
- 索引构建:O(n)
-
层级聚合:O(n)
-
内存占用测试 :
- 100 万节点:约 500MB 堆内存
-
通过对象复用可降低 30% 内存
-
基准测试对比 :
- 相比原生 SQL 方案,处理速度提升 5 - 8 倍
- 内存占用减少 40%
生产环境建议
- 批量处理大小 :
- 建议每批次处理 10-50 万条记录
-
过大会导致 GC 压力,过小影响吞吐量
-
JVM 调优 :
- 设置合理的年轻代大小 (-Xmn)
- 使用 G1 垃圾收集器
-
增大直接内存 (-XX:MaxDirectMemorySize)
-
监控指标 :
- 处理耗时百分位 (TP99)
- 内存使用趋势
- GC 频率和耗时
延伸思考
这个方案可以扩展到其他 OLAP 场景,如:
- 时间序列数据的层级聚合
- 多维度交叉分析
- 实时数据预处理
但同时也留下了几个值得思考的问题:
- 如何进一步降低内存占用?
- 能否利用 ClickHouse 本身的聚合能力减少数据传输?
- 对于超大规模数据,如何实现分布式处理?
这些问题的解决将帮助我们构建更高效的数据处理管道。
正文完
发表至: 技术分享
近一天内
