所有的代码分析基于 polars 1.43.0 版本
分析polars的Expressions and contexts
如果说 Series 是 Polars 的数据单元、DataFrame 是组织单元,那么 表达式 (Expression) 就是 Polars 的”指令单元”:它是一个惰性的数据变换描述,不执行任何计算;而 上下文 (Context) 是它的执行环境——同一个表达式在不同 context 里会产生不同结果。这两者构成了 Polars 的 DSL(领域专用语言)核心。
架构总览
1. 表达式:惰性 DSL
1.1 Python 层:Expr 类的构建
Python 的 Expr 类定义在:
class Expr:
"""Expressions that can be used in various contexts."""
# NOTE: This `= None` is needed to generate the docs with sphinx_accessor.
_pyexpr: PyExpr = None # type: ignore[assignment]展开折叠代码 (140-151 行,共 12 行)
_accessors: ClassVar[set[str_]] = {
"arr",
"bin",
"cat",
"dt",
"ext",
"list",
"meta",
"name",
"str",
"struct",
}
@property
def bin(self) -> ExprBinaryNameSpace:
"""
Create an object namespace of all binary related methods.
See the individual method pages for full details
"""
return ExprBinaryNameSpace(self)关键设计:
_pyexpr: PyExpr— 和 Series/DataFrame 一样,Python 层不存数据也不存逻辑,只是 Rust 侧PyExpr的薄壳_accessors— 声明arr/dt/list/meta/name/str/struct等命名空间,每个@property返回对应的 NameSpace 类,把同名操作按数据类型分组- 运算符重载 —
bmi_expr = pl.col("weight") / (pl.col("height") ** 2)之所以能这么写,是因为Expr.__truediv__、__pow__等被重载,每个运算符都构建一个新的PyExpr
运算符的实现模式:
def __add__(self, other: IntoExpr) -> Expr:
other_pyexpr = parse_into_expression(other, str_as_lit=True)
return wrap_expr(self._pyexpr + other_pyexpr)展开折叠代码 (342-353 行,共 12 行)
def __radd__(self, other: IntoExpr) -> Expr:
other_pyexpr = parse_into_expression(other, str_as_lit=True)
return wrap_expr(other_pyexpr + self._pyexpr)
def __and__(self, other: IntoExprColumn | int | bool) -> Expr:
other_pyexpr = parse_into_expression(other)
return wrap_expr(self._pyexpr.and_(other_pyexpr))
def __rand__(self, other: IntoExprColumn | int | bool) -> Expr:
other_expr = parse_into_expression(other)
return wrap_expr(other_expr.and_(self._pyexpr)) def __eq__(self, other: IntoExpr) -> Expr: # type: ignore[override]
warn_null_comparison(other)
other_pyexpr = parse_into_expression(other, str_as_lit=True)
return wrap_expr(self._pyexpr.eq(other_pyexpr))
def __floordiv__(self, other: IntoExpr) -> Expr:
other_pyexpr = parse_into_expression(other)
return wrap_expr(self._pyexpr // other_pyexpr)所有输入(字符串、数字、列表、Series、Expr)都先经过 parse_into_expression 统一转成 PyExpr,再交给 Rust 端做二元运算。__eq__ 被重载成表达式比较(==),所以 pl.col("a") == 1 返回的是布尔表达式而不是 Python bool——同时 __bool__ 被禁用并给出改写法提示(expr.py:321)。
parse_into_expression 的类型归一化:
def parse_into_expression(
input: IntoExpr,
*,
str_as_lit: bool = False,
list_as_series: bool = False,
structify: bool = False,
dtype: PolarsDataType | None = None,
require_selector: bool = False,
) -> PyExpr:展开折叠代码 (31-55 行,共 25 行)
"""
Parse a single input into an expression.
Parameters
----------
input
The input to be parsed as an expression.
str_as_lit
Interpret string input as a string literal. If set to `False` (default),
strings are parsed as column names.
list_as_series
Interpret list input as a Series literal. If set to `False` (default),
lists are parsed as list literals.
structify
Convert multi-column expressions to a single struct expression.
dtype
If the input is expected to resolve to a literal with a known dtype, pass
this to the `lit` constructor.
require_selector
Require that the input is a valid selector (eg: column name or selector).
Returns
-------
PyExpr
""" if isinstance(input, pl.Expr):
expr = input
if structify:
expr = _structify_expression(expr)
elif isinstance(input, str) and not str_as_lit:
expr = F.col(input)
else:
if require_selector:
msg = f"cannot turn {qualified_type_name(input)!r} into selector"
raise TypeError(msg)
elif isinstance(input, list) and list_as_series:
expr = F.lit(pl.Series(input), dtype=dtype)
else:
expr = F.lit(input, dtype=dtype)核心规则:字符串默认当列名(F.col(input)),其他类型当字面量(F.lit(input))。str_as_lit 参数控制字符串是列名还是字符串常量——这就是 df.select("a")(字符串 “a” 被解析为列名)与 pl.col("a") == "x"(运算符重载里传入 str_as_lit=True,“x” 是字符串字面量)语义不同的根源。
1.2 桥接层:PyExpr 透明包装
#[pyclass(from_py_object)] // Not marked as frozen for pickling, but that's the only &mut self method.
#[repr(transparent)]
#[derive(Clone)]
pub struct PyExpr {
pub inner: Expr,
}
impl From<Expr> for PyExpr {
fn from(expr: Expr) -> Self {
PyExpr { inner: expr }
}
}PyExpr { inner: Expr } 用 #[repr(transparent)] 保证内存布局与 Expr 完全一致,因此 Vec<PyExpr> 和 Vec<Expr> 之间可以零成本 transmute 互转(ToExprs trait 直接用 Vec::from_raw_parts 把指针 reinterpret,见 expr/mod.rs:54-63)。
pl.col("weight") 在 Rust 侧只是构造 dsl::Expr::Column:
#[pyfunction]
pub fn coalesce(exprs: Vec<PyExpr>) -> PyExpr {
let exprs = exprs.to_exprs();
dsl::coalesce(&exprs).into()
}
#[pyfunction]
pub fn col(name: &str) -> PyExpr {
dsl::col(name).into()
}
#[pyfunction]
pub fn element() -> PyExpr {
dsl::element().into()
}1.3 Rust DSL:Expr 枚举(AST)
表达式真正落地为 Expr 枚举——一个递归的 AST:
/// Expressions that can be used in various contexts.
///
/// Queries consist of multiple expressions.
/// When using the polars lazy API, don't construct an `Expr` directly; instead, create one using
/// the functions in the `polars_lazy::dsl` module. See that module's docs for more info.
#[derive(Clone, PartialEq)]
#[must_use]
#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
#[cfg_attr(feature = "dsl-schema", derive(schemars::JsonSchema))]
pub enum Expr {
/// Values in a `eval` context.
///
/// Equivalent of `pl.element()`.
Element,
Alias(Arc<Expr>, PlSmallStr),
Column(PlSmallStr),
Selector(Selector),
Literal(LiteralValue),
DataTypeFunction(DataTypeFunction),
BinaryExpr {展开折叠代码 (77-80 行,共 4 行)
left: Arc<Expr>,
op: Operator,
right: Arc<Expr>,
}, Cast {
expr: Arc<Expr>,
dtype: DataTypeExpr,
options: CastOptions,
},
Sort {
expr: Arc<Expr>,
options: SortOptions,
},
Gather {
expr: Arc<Expr>,
idx: Arc<Expr>,
returns_scalar: bool,
null_on_oob: bool,
},
SortBy {
expr: Arc<Expr>,
by: Vec<Expr>,
sort_options: SortMultipleOptions,
},
Agg(AggExpr),
/// A ternary operation
/// if true then "foo" else "bar"
Ternary {
predicate: Arc<Expr>,
truthy: Arc<Expr>,
falsy: Arc<Expr>,
},
Function {
/// function arguments
input: Vec<Expr>,
/// function to apply
function: FunctionExpr,
},
Explode {| 变体 | 含义 |
|---|---|
Column(PlSmallStr) | 引用一列 |
Literal(LiteralValue) | 字面量(数字、字符串、空值…) |
BinaryExpr { left, op, right } | 二元运算,Arc<Expr> 递归 |
Agg(AggExpr) | 聚合(sum/mean/implode…) |
Ternary { predicate, truthy, falsy } | 三元 when/then/otherwise |
Function { input, function } | 命名函数,FunctionExpr 是内置函数注册表 |
Alias(Arc<Expr>, PlSmallStr) | 重命名 |
Filter / Sort / Cast / Slice / Over… | 其他变换 |
设计亮点:
- 子表达式用
Arc<Expr>而非Box<Expr>— 同一表达式子树(如bmi_expr被复用 3 次)可共享,为公共子表达式消除(CSE)提供基础 Expr本身就是 AST,不是闭包 — 因此可以被打印(__repr__输出[(col("weight")) / (col("height").pow(...))])、被遍历、被重写、被序列化(支持 serde + cloud 远程执行)FunctionExpr是穷举的枚举而非函数指针 — 优化器可以对每个函数做静态分析(是否确定性、是否可并行、是否输入可展开),见下文 expression expansion
聚合是独立的 AggExpr 枚举(expr/mod.rs:22-55),因为聚合的语义(对一组值归约)和逐元素操作完全不同,物理执行时走不同的求值路径。
2. Contexts:四种执行环境
Polars 文档定义了 4 个 context,它们在逻辑计划 IR 里各有对应节点:
| Context | Python 方法 | IR 节点 | 语义 |
|---|---|---|---|
| select | df.select | IR::Select | 只保留表达式选出的列 |
| with_columns | df.with_columns | IR::HStack | 原表 + 新增列 |
| filter | df.filter | IR::Filter | 按布尔表达式筛行 |
| group_by + agg | df.group_by(...).agg(...) | IR::GroupBy | 分组后聚合 |
// Polars' `select` operation. This may access full materialized data.
Select {
input: Node,
expr: Vec<ExprIR>,
schema: SchemaRef,
options: ProjectionOptions,
},展开折叠代码 (91-100 行,共 10 行)
Sort {
input: Node,
by_column: Vec<ExprIR>,
slice: Option<(i64, usize, Option<DynamicPred>)>,
sort_options: SortMultipleOptions,
},
Cache {
input: Node,
/// This holds the `Arc<DslPlan>` to guarantee uniqueness.
id: UniqueId, },
GroupBy {
input: Node,
keys: Vec<ExprIR>,
aggs: Vec<ExprIR>,
schema: SchemaRef,
maintain_order: bool,
options: Arc<GroupbyOptions>,
apply: Option<PlanCallback<DataFrame, DataFrame>>,
},展开折叠代码 (111-123 行,共 13 行)
Join {
input_left: Node,
input_right: Node,
schema: SchemaRef,
left_on: Vec<ExprIR>,
right_on: Vec<ExprIR>,
options: Arc<JoinOptionsIR>,
},
Gather {
input: Node,
idxs: Node,
null_on_oob: bool,
}, HStack {
input: Node,
exprs: Vec<ExprIR>,
schema: SchemaRef,
options: ProjectionOptions,
},
Distinct {注意 IR::HStack 的注释(“Horizontal stack”)——with_columns 本质是横向堆叠新列到原表,所以它的表达式必须产出与行数等长的 Series;而 Select 允许标量(文档说 select 里的标量会被 broadcast 到与最长列等长)。
Python 侧四个 context 全部收敛到同一个 lazy 引擎(eager 只是立刻 collect,见第 04 篇),因此 expression expansion、谓词下推等优化对 eager/lazy 一视同仁。
3. 物理执行:Executor 树
3.1 Executor trait
// Executor are the executors of the physical plan and produce DataFrames. They
// combine physical expressions, which produce Series.
/// Executors will evaluate physical expressions and collect them in a DataFrame.
///
/// Executors have other executors as input. By having a tree of executors we can execute the
/// physical plan until the last executor is evaluated.
pub trait Executor: Send + Sync {
fn execute(&mut self, cache: &mut ExecutionState) -> PolarsResult<DataFrame>;
fn is_cache_prefiller(&self) -> bool {
false
}
}collect(crates/polars-lazy/src/frame/mod.rs:746)最终把优化后的 IR 转成物理计划:一棵 Executor 树。每个 Executor 从输入 Executor 拿到 DataFrame,评估物理表达式得到新列,输出新 DataFrame。四个 context 对应四个 Executor:
| Context | Executor | 文件 |
|---|---|---|
| select | ProjectionExec | executors/projection.rs |
| with_columns | StackExec | executors/stack.rs |
| filter | FilterExec | executors/filter.rs |
| group_by | GroupByExec | executors/group_by.rs |
3.2 PhysicalExpr trait:表达式求值
DSL 的 Expr(AST)在物理计划阶段被编译成 Arc<dyn PhysicalExpr>——一个 trait object,屏蔽具体类型:
pub trait PhysicalExpr: Send + Sync {
fn as_expression(&self) -> Option<&Expr> {
None
}
fn as_column(&self) -> Option<PlSmallStr> {
None
}
/// Take a DataFrame and evaluate the expression.
///
/// Note: implementers should implement evaluate_impl instead, as this wraps
/// that call with an error context.
fn evaluate(&self, df: &DataFrame, state: &ExecutionState) -> PolarsResult<Column> {
self.evaluate_impl(df, state).map_err(|e| {
if let Some(expr) = self.as_expression() {
e.with_expr_context(expr.to_string().into())
} else {
e
}
})
}
fn evaluate_impl(&self, df: &DataFrame, _state: &ExecutionState) -> PolarsResult<Column>;这个 trait 有两个关键方法,对应两种 context:
evaluate(df, state)— 在整表上求值,用于 select/with_columns/filterevaluate_on_groups(df, groups, state)— 在分组索引上求值,用于 group_by 的聚合表达式
3.3 表达式的并行求值
select/with_columns 的表达式之间是”令人尴尬的并行”(embarrassingly parallel),直接 rayon 并行:
fn run_exprs_par(
df: &DataFrame,
exprs: &[Arc<dyn PhysicalExpr>],
state: &ExecutionState,
) -> PolarsResult<Vec<Column>> {
RAYON.install(|| {
exprs
.par_iter()
.map(|expr| expr.evaluate(df, state))
.collect()
})
}
fn run_exprs_seq(
df: &DataFrame,
exprs: &[Arc<dyn PhysicalExpr>],
state: &ExecutionState,
) -> PolarsResult<Vec<Column>> {
exprs.iter().map(|expr| expr.evaluate(df, state)).collect()
}
pub(super) fn evaluate_physical_expressions(
df: &mut DataFrame,
exprs: &[Arc<dyn PhysicalExpr>],
state: &ExecutionState,
has_windows: bool,
run_parallel: bool,
) -> PolarsResult<Vec<Column>> {
let expr_runner = if has_windows {
execute_projection_cached_window_fns
} else if run_parallel && exprs.len() > 1 {
run_exprs_par
} else {
run_exprs_seq
};
let selected_columns = expr_runner(df, exprs, state)?;
if has_windows {
state.clear_window_expr_cache();
}
Ok(selected_columns)
}调度决策在 evaluate_physical_expressions:
- 有窗口函数 →
execute_projection_cached_window_fns(需要缓存 + 串行,因为窗口函数相互依赖) - 无窗口且表达式 > 1 个 →
run_exprs_par(par_iter并行求值每个表达式) - 否则 →
run_exprs_seq
ProjectionExec 甚至还有垂直并行:当 DataFrame 多 chunk 且足够高时,按 chunk 切分后 par_iter 并行执行整条 select,最后纵向拼接(projection.rs:26-51)。FilterExec 也有同样的分块并行逻辑。
4. 为什么需要 AST 和 Executor 两层?
上面出现了 Expr(AST)和 Executor(物理执行)两套结构,这是初学者最容易困惑的地方:有了 AST,不就知道怎么执行了吗? 其实不是——完整的链路是三层:
Expr (DSL, 你写的) → AExpr + IR (逻辑计划) → PhysicalExpr + Executor (物理计划)
| 层 | 回答的问题 | 谁在用 |
|---|---|---|
Expr / AExpr | 算什么? | 优化器重写它 |
IR | 数据怎么流? | 谓词下推、投影剪枝 |
PhysicalExpr + Executor | 具体怎么跑? | CPU 实际执行 |
核心区别一句话:AST 回答”做什么”,Executor 回答”怎么做”。
4.1 AST 是声明式,Executor 是命令式
pl.col("weight") / (pl.col("height") ** 2) 只声明了”要做除法”,没指定:
- 按列遍历还是按 chunk 分块并行?
- 中间值放哪块内存、要不要缓存?
- 结果以什么形式产出(整列 / 标量 / 分组聚合)?
这些执行细节 AST 一概没有,需要在编译成物理计划时根据数据规模、线程数、上下文决定。
4.2 AST 必须能被优化器重写
谓词下推、投影剪枝、公共子表达式消除(CSE)、join 重排,本质都是对 AST/逻辑计划做等价变换。如果 AST 一开始就绑定了执行方式,优化器就没有可改的空间了。AST 必须保持”纯净的数学语义”,才能被自由重排。
4.3 一份 AST 可编译到多个执行器
同一个逻辑计划可以跑在:
- 内存引擎
polars-mem-engine(本系列分析的这个) - 流式引擎(分批执行,控制内存峰值)
- Cloud 分布式执行(把 AExpr 序列化发送到远端节点)
- GPU 引擎
如果 AST 定义了执行,每加一个引擎就得改所有表达式。解耦之后,新增引擎只是多写一个”编译器”(把 IR 编译成该引擎的执行器),AST 完全不动。
4.4 编译阶段消解 AST 的多义性
AST → 物理表达式的过程是真正的”编译”,需要 schema 和上下文信息(AST 里没有):
pl.col("weight", "height")展开成两个独立表达式pl.col(pl.Float64)按输入 schema 展开成 N 个- 标量列要不要广播
- 聚合走
evaluate(整表)还是evaluate_on_groups(分组)
同一个 pl.col("x").mean(),在 select 里编译成整表聚合,在 group_by().agg 里编译成分组聚合——AST 完全一样,编译结果不同,因为编译时携带的上下文不同。
4.5 执行期开销
Expr是Arc<Expr>递归 + 模式匹配分派,还要遍历 arenaPhysicalExpr是扁平的Arc<dyn PhysicalExpr>,一个evaluate(&df)直接返回Column,没有 AST 遍历
物理表达式在规划时已经完成了 to_field 类型推导(输出 schema 提前确定),执行时零决策成本。
类比:SQL 的
SELECT a, SUM(b) FROM t GROUP BY a是 AST,PostgreSQL / DuckDB / Spark 各自把它编译成完全不同的执行计划。理解”这是分组求和”是一回事(AST),“怎么算最快——哈希分组还是排序分组、要不要并行、怎么下推”是另一回事(物理执行)。
5. 四种 Context 的差异本质
5.1 select 与 with_columns:长度语义
两者共享 evaluate_physical_expressions,差别在收尾的 check_expand_literals:
- select — 允许长度 1 的标量列,按最长列
broadcast展开(projection_utils.rs:347-363),所以ideal_max_bmi=25能被广播成整列 - with_columns —
StackExec用df.with_columns_mut(res, schema)要求新列与行数一致(stack.rs:44)
5.2 filter:布尔掩码
FilterExec 先求值谓词表达式,再把它转成布尔掩码:
fn execute_hor(
&mut self,
df: DataFrame,
state: &mut ExecutionState,
) -> PolarsResult<DataFrame> {
if self.has_window {
state.insert_has_window_function_flag()
}
let c = self.predicate.evaluate(&df, state)?;
if self.has_window {
state.clear_window_expr_cache()
}
// @scalar-opt
// @partition-opt
df.filter(column_to_mask(&c, df.height())?.as_ref())
}column_to_mask(filter.rs:13-43)专门处理 Column::Scalar——布尔常量列不用物化成整列,直接取 Scalar 值构造长度 1 的掩码,交给 DataFrame 的列级 filter。
5.3 group_by:聚合表达式在分组上求值
GroupByExec 先算分组键,得到 GroupPositions(每个组的下标集合),然后聚合表达式通过 evaluate_on_groups 在分组上求值:
pub(super) fn evaluate_aggs(
df: &DataFrame,
aggs: &[Arc<dyn PhysicalExpr>],
groups: &GroupPositions,
state: &ExecutionState,
) -> PolarsResult<Vec<Column>> {
RAYON.install(|| {
aggs.par_iter()
.map(|expr| {
let agg = expr.evaluate_on_groups(df, groups, state)?.finalize();
polars_ensure!(agg.len() == groups.len(), agg_len = agg.len(), groups.len());
Ok(agg)
})
.collect::<PolarsResult<Vec<_>>>()
})
}每个聚合表达式并行执行 evaluate_on_groups(df, groups, state),返回长度等于组数的列。这就是”同一表达式在不同 context 产生不同结果”的物理根源:pl.col("name").mean() 在 select 里是整列求均值(1 个值),在 group_by.agg 里是每组求均值(N 个值)——区别只在于求值入口走 evaluate 还是 evaluate_on_groups。
6. Expression Expansion:表达式的列级展开
pl.col("weight", "height").mean() 一个表达式会变成两列输出。这个展开发生在 DSL 转 IR 阶段(rewrite_projections / expand_expression):
//! this contains code used for rewriting projections, expanding wildcards, regex selection etc.
use super::*;
use crate::constants::{
POLARS_ELEMENT, POLARS_STRUCTFIELDS, get_pl_element_name, get_pl_structfields_name,
};
pub fn prepare_projection(
exprs: Vec<Expr>,
schema: &Schema,
opt_flags: &mut OptFlags,
) -> PolarsResult<(Vec<Expr>, Schema)> {
let exprs = rewrite_projections(exprs, &PlIndexSet::new(), schema, opt_flags)?;
let schema = expressions_to_schema(&exprs, schema, |duplicate_name: &str| {
format!("projections contained duplicate output name '{duplicate_name}'")
})?;
Ok((exprs, schema))
}展开折叠代码 (20-22 行,共 3 行)
pub fn is_regex_projection(name: &str) -> bool {
name.starts_with('^') && name.ends_with('$')
}
pub fn expand_expression(
expr: &Expr,
ignored_selector_columns: &PlIndexSet<PlSmallStr>,
schema: &Schema,
out: &mut Vec<Expr>,
opt_flags: &mut OptFlags,
) -> PolarsResult<()> {
if expr.into_iter().all(|e| !needs_expansion(e)) {
out.push(expr.clone());
return Ok(());
}
expand_expression_rec(expr, ignored_selector_columns, schema, out, opt_flags)?;
Ok(())
}展开折叠代码 (40-59 行,共 20 行)
/// In case of single col(*) -> do nothing, no selection is the same as select all
/// In other cases replace the wildcard with an expression with all columns
pub fn rewrite_projections(
exprs: Vec<Expr>,
ignored_selector_columns: &PlIndexSet<PlSmallStr>,
schema: &Schema,
opt_flags: &mut OptFlags,
) -> PolarsResult<Vec<Expr>> {
let mut result = Vec::with_capacity(exprs.len() + schema.len());
for expr in &exprs {
expand_expression(
expr,
ignored_selector_columns,
schema,
&mut result,
opt_flags,
)?;
}
Ok(result)}设计亮点:
needs_expansion先做整棵树检查(expr_expansion.rs:31)— 无 Selector/通配符的表达式直接原样入队,零开销- 展开发生在 IR 构建期,而非执行期 — 展开后每条表达式是独立的物理表达式,天然可以被
run_exprs_par并行 - 依赖 schema —
(pl.col(pl.Float64) * 1.1)展开成几列取决于输入 schema,所以文档说”无法预先知道会展开成多少个表达式”,这正是rewrite_projections(exprs, ..., schema, ...)需要 schema 参数的原因 - 特殊的函数展开规则 —
function_input_wildcard_expansion(expr_expansion.rs:77)枚举哪些函数(如concat_str、max_horizontal、coalesce)会把通配符展开进输入而非自身,语义精细到每个函数
设计亮点总结
-
表达式即数据(AST 而非闭包) —
Expr是递归枚举,可打印、可遍历、可重写、可序列化(serde + cloud)。惰性不只是”不执行”,而是让整个查询成为可分析、可优化的对象 -
Arc<Expr>共享子树 + CSE — 子表达式用Arc而非Box,复用bmi_expr三次时底层共享同一 AST 子树,优化器可以消除重复计算 -
运算符重载构建 DSL — Python 通过
__truediv__/__pow__/__eq__重载让pl.col("weight") / (pl.col("height") ** 2)这种数学式写法成立;所有输入经parse_into_expression统一归一化为 PyExpr -
零成本桥接 —
#[repr(transparent)]让PyExpr与Expr内存布局一致,Vec<PyExpr> ↔ Vec<Expr>直接指针 reinterpret,无拷贝 -
逻辑 / 物理分层 — AST 声明”做什么”,Executor 声明”怎么做”:AST 保持纯净语义供优化器自由重写,一份逻辑计划可编译到内存 / 流式 / Cloud / GPU 多引擎;编译期消解多义性(展开、广播、求值入口),物理表达式执行期零 AST 开销
-
Context 决定求值路径 — 同一表达式在 select 走
evaluate(整表),在 group_by.agg 走evaluate_on_groups(按组),物理差异只在求值入口;PhysicalExpr的这两个方法就是 context 的落点 -
上下文无关的并行 — 表达式间并行(
run_exprs_par)对所有 context 统一;窗口函数识别后自动退回缓存串行路径;分 chunk 的垂直并行让 select/filter 也享受并行 -
展开期优化 — expression expansion 在 DSL→IR 阶段按 schema 把 Selector/通配符展开成独立表达式,展开结果天然并行,且展开规则精确到每个内置函数