跳到主要内容
版本:1.0

分布式查询

Frontend 和 Datanode 使用同一套基于 DataFusion 的查询引擎。在分布式模式下,Frontend 会增加一个规划步骤,将 Datanode 上执行的工作与 Frontend 上完成的工作分开。

Frontend query

分布式规划​

分布式规划器重写逻辑计划,把可以下推的算子移向表扫描,并用 MergeScan 节点包装远端子计划。分区列上的谓词还会在任务调度前用于裁剪 Region。

算子能否下推取决于计划形态和算子本身的性质。不支持的部分会保留在 Frontend。初始设计及交换律规则参见分布式规划器 RFC。

分布式计划​

远端输入是完整的逻辑子计划,并不局限于表扫描。Frontend 使用 Substrait 序列化子计划,再向持有相应数据的 Datanode 发送 Region 级请求。Datanode 在本地规划并执行子计划,将结果流返回 Frontend。Frontend 合并远端数据流,并执行没有下推的算子。