一、问题背景:跑数需求为何成为工程师的"噩梦"
在大型互联网公司中,业务方对数据的需求是海量的、多样的、且频繁变更的。产品经理需要看转化漏斗数据,运营需要查看用户分层报表,技术负责人需要监控系统的关键指标。每一类需求最终都归结为"写个SQL把数据拉出来"。
在平台化之前,这类跑数需求的基本流程是:业务方提需求 -> 数据工程师写SQL -> 排队等资源 -> 执行 -> 出结果。一个简单的需求从提出到拿到数据,平均需要1-2个工作日。而在高峰期,积压的跑数需求可以达到上百个,数据工程师的大部分时间都花在了重复性的SQL编写和结果交付上。
更棘手的问题是跨库查询。百度的业务数据分散在MySQL、HBase、Phoenix(HBase上的SQL层)、ClickHouse等多个存储系统中,一个看似简单的查询可能需要关联三四个不同数据源的数据。传统的做法是先通过ETL把数据同步到同一个数据仓库中再查询,但这引入了数据延迟,对于需要实时或准实时数据的需求来说无法接受。
我们需要一个系统,让业务方能够自助地通过简单的界面完成大部分跑数需求,同时底层能够透明地处理多数据源、跨库关联等复杂逻辑。这就是数据自助查询平台的出发点。
二、管道设计思想:算子驱动的数据处理流水线
平台的核心设计思想借鉴了编译原理中的抽象语法树(AST)概念和大数据处理中的算子树模型。一个SQL查询不再是"一次性执行完"的操作,而是被解析为一棵算子树,每个节点代表一个数据处理算子,数据像水流一样从叶子节点流经中间算子,最终在根节点汇聚输出。
算子系统的设计是我们整个平台的基石。我们定义了以下核心算子类型:
- DataSourceOperator:数据源读取算子,负责从MySQL、HBase、Phoenix等不同数据源拉取数据,内部封装了各数据源的连接池管理和分页读取策略。
- FilterOperator:过滤算子,支持WHERE条件下推,尽可能在数据源层面完成过滤,减少网络传输量。
- JoinOperator:关联算子,支持Hash Join和Nested Loop Join两种策略,并能够根据数据量自动选择最优策略。
- AggregateOperator:聚合算子,处理GROUP BY、COUNT、SUM、AVG等聚合操作,支持内存聚合和溢写磁盘两种模式。
- SortOperator:排序算子,支持ORDER BY和LIMIT,实现了外部归并排序以处理超大数据集。
- ProjectOperator:投影算子,负责SELECT字段的裁剪,减少后续算子需要处理的数据量。
每个算子都实现了统一的 PipelineOperator 接口,通过 push() 和 pull() 方法进行数据流的驱动:
public interface PipelineOperator {
/**
* 初始化算子,解析算子特定的配置参数
*/
void init(OperatorConfig config);
/**
* 推模式:上游算子将数据行推送到当前算子
* 返回 true 表示继续接收,false 表示当前算子已达到处理上限
*/
boolean push(RowData row);
/**
* 拉模式:下游算子从当前算子拉取数据
*/
List<RowData> pull(int batchSize);
/**
* 流式计算完成后的收尾操作(如输出缓冲区剩余数据)
*/
void finish();
/**
* 释放资源
*/
void close();
}
流式计算的优势在于内存效率。传统的批量处理模型需要将所有数据加载到内存后再进行下一步处理,而流式模型允许数据逐行流经算子链,中间结果无需全量驻留内存。对于大结果集查询,这种模型可以将内存占用降低一个数量级。
三、跨库Join与自动Join:让复杂查询透明化
跨库Join是平台最具技术挑战的部分。业务方只需要写一个标准的SQL,系统需要自动判断哪些表在哪个数据源,然后生成最优的执行计划。我们的实现分为三个层次:
第一层:元数据驱动的数据源路由。我们构建了一个统一的元数据中心,记录了每张表的物理位置、Schema信息、索引情况、数据量级等元数据。当用户提交SQL后,系统首先通过元数据中心确定每个表的数据源:
public class MetadataRouter {
private final MetadataRepository metadataRepo;
/**
* 解析SQL中涉及的表,并路由到对应的数据源
*/
public ExecutionPlan route(String sql) {
List<TableRef> tables = SqlParser.extractTables(sql);
Map<TableRef, DataSourceInfo> routing = new HashMap<>();
for (TableRef table : tables) {
DataSourceInfo ds = metadataRepo
.findDataSource(table.getDatabase(), table.getName());
if (ds == null) {
throw new RouteException("表 " + table
+ " 未在元数据中心注册");
}
routing.put(table, ds);
}
return ExecutionPlan.build(routing, sql);
}
}
第二层:基于代价的Join策略选择。当两个表分布在不同数据源时,系统需要决定将哪个表的数据拉到本地来执行Join。我们的策略是:选择数据量较小的表(Build Table),将其全量拉到内存中构建Hash表,然后逐行拉取较大表(Probe Table)的数据进行Hash匹配。当Build Table超过内存阈值时,自动降级为Nested Loop Join + 分批处理。
第三层:自动Join推断。很多业务方并不清楚表之间的关联关系。我们通过分析表的外键关系、命名规范(如 user_id、order_id 等公共字段)、历史查询日志,自动推断表之间的Join条件。当用户查询了用户表和订单表但没有写Join条件时,系统会提示:"检测到两张表可能需要通过 user_id 关联,是否自动补充Join条件?"
public class AutoJoinInferencer {
/**
* 基于公共字段推断Join条件
*/
public JoinCondition infer(List<TableSchema> tables) {
// 1. 收集所有表的字段信息
List<ColumnStat> allColumns = tables.stream()
.flatMap(t -> t.getColumns().stream())
.collect(Collectors.toList());
// 2. 按字段名分组,寻找跨表的同名字段
Map<String, List<ColumnStat>> byName = allColumns.stream()
.collect(Collectors.groupingBy(ColumnStat::getName));
// 3. 筛选出出现在多个表中的字段作为候选Join Key
List<JoinCandidate> candidates = byName.entrySet().stream()
.filter(e -> e.getValue().size() > 1)
.map(e -> new JoinCandidate(
e.getKey(), e.getValue()))
.sorted(Comparator.comparingDouble(
JoinCandidate::getConfidence).reversed())
.collect(Collectors.toList());
return candidates.isEmpty() ? null
: candidates.get(0).toJoinCondition();
}
}
通过这套跨库Join和自动Join机制,业务方无需关心数据在哪个库、表之间如何关联,只需要用熟悉的SQL表达业务意图即可。上线后,平台日均处理查询请求超过5000次,80%以上的跑数需求实现了自助完成,数据工程师的重复性工作量下降了约70%,释放出的精力被投入到数据质量治理和更深层次的分析工作中。
回顾这个项目,我认为最重要的设计决策是将复杂度封装在平台内部,对外暴露最简单的接口。管道式架构让系统具备了良好的可扩展性——新增一个数据源只需要实现对应的 DataSourceOperator,新增一种数据处理逻辑只需要增加一个算子类型,而不会影响整体架构的稳定性。