在 grouped_sum 的外层 group 循环中,对未来第 4 个 group 的 indices 头部地址做软件预取:
for (gi, g) in groups.iter().enumerate() {
if gi + 4 < n_groups {
unsafe {
sve_sum::prefetch_next_group(&groups[gi + 4]); // prfm pldl1keep
}
}
// process current group via SVE kernel
}
pldl1keep(temporal hint)用于 indices 数据,因为它会在内层 SVE 循环中被重复读取,应保留在 L1。value 数据数组的预取在 SVE 内核内部由 pldl1strm(streaming hint)处理。
letmut any = false;
let (mut a0, mut a1, mut a2, mut a3) = (0, 0, 0, 0);
while i + 4 <= n {
if nulls.is_valid(i0) { a0 = a0.wrapping_add(*values.add(i0)); any = true; }
if nulls.is_valid(i1) { a1 = a1.wrapping_add(*values.add(i1)); any = true; }
if nulls.is_valid(i2) { a2 = a2.wrapping_add(*values.add(i2)); any = true; }
if nulls.is_valid(i3) { a3 = a3.wrapping_add(*values.add(i3)); any = true; }
i += 4;
}
if any { Some(a0.wrapping_add(a1).wrapping_add(a2).wrapping_add(a3)) } else { None }
Is your feature request related to a problem?
no
Describe the solution you'd like
1 前言
1.1 背景
在 aarch64 和 X86 平台进行基准测试,测试结果显示 Daft 的部分场景下在 aarch64 架构上的性能相比 X86 劣化。Daft 在一些场景下,依赖了
daft_core::array::ops::sum::grouped_sum接口且该接口性能在不同的场景下均有性能劣化,本文主要分析 daft 中grouped_sum函数性能劣化原因和优化方法。2 Story概述
2.1 Story需求描述(必填)
2.2 Story功能描述(必填)
grouped_sum函数是 Daft 聚合算子中的核心分组求和函数,负责对按GroupIndices(Vec<Vec<u64>>)描述的每个分组内的若干行值进行求和。该函数通过impl_daft_numeric_agg!宏统一为 5 种数值类型生成实现:Int64Type/UInt64Type/Float32Type/Float64Type/Decimal128Type。该函数的核心工作流程包括:
as_arrow()获取底层 ArrowPrimitiveArraynull_count() > 0选择无 null 快速路径或含 null 慢速路径g.iter().fold(...)将该 group 中所有索引对应的值累加Vec::from_iter收集为输出DataArray原始实现中存在以下性能瓶颈:
Vec<Vec<u64>>pointer chase:跨 group 切换时硬件预取器失效,产生 L1 cache missas_arrow()重复 downcast:每个self.get(idx)内部都做一次Any::downcast_ref,在大数据量下累积可观开销2.3 Story用户使用场景分析(必填)
一、用户场景分析
groupby(...).sum()或agg(col(...).sum())操作时,grouped_sum是核心路径。grouped_sum的性能直接影响整体查询性能。Q1 是典型的 sum-heavy group-by 查询,对 lineitem 表按l_returnflag和l_linestatus分组后对多列做求和。二、新增/变更的文件、脚本 (必填)
本次接口优化需要修改的文件如下:
2.4 Story约束(必填)
基于 Daft v0.7.5 版本进行优化。
3 Story设计描述
3.1 Story设计思路(必填)
问题分析
原始实现中
grouped_sum存在以下性能瓶颈:fold 单累加器串行依赖:
g.iter().fold(0, |acc, idx| acc + self.get(idx).unwrap())形成 loop-carried dependency chain,下一次add必须等上一次完成。CPU 无法并行 issue 多个独立的 add 指令,IPC 严重受限于 add latency。f64 的faddlatency 在鲲鹏 920B 上是 4-5 周期,单累加器把 IPC 压到 ~0.2。Vec<Vec<u64>>两级 Vec pointer chase:外层 Vec 顺序访问 OK,但每个 innerVec<u64>是独立堆分配,group 之间存在 pointer chase。鲲鹏 920B 的硬件预取器对跨 group 的随机布局无效,反映在 perf stat 上是 920B cache-misses(2.71B)比 9654(1.92B)高 41%。as_arrow()在get(idx)中重复 downcast:get内部每次调用都做Any::downcast_ref+Result::unwrap,在 group size 为数千万行的场景下累计上亿次冗余 downcast。未利用 SVE 指令集:鲲鹏 920B 拥有 SVE 1.0(VL=128),其中
ld1d {z1.d}, p0/z, [base, z0.d, lsl #3](向量 gather)和faddv(FP 横向归约)是 x86 没有高效等价物的指令。LLVM auto-vectorizer 无法为 fold 模式生成这些指令,需要手写内联汇编。null 路径
Option模式匹配开销:含 null 时 fold 中每次迭代的match (acc, self.get(idx))都要做 4-way pattern match,分支预测压力大。优化设计
优化一:SVE gather + 谓词累加 + uaddv 横向归约(i64 / u64)
为 i64 / u64 实现 SVE 内核,使用一条
ld1d指令一次性 gather 整个 SVE 向量长度的值(920B VL=128 时一次 2 个 i64),用谓词受控的add累加,最后用uaddv做横向归约。优化前:
g.iter().fold(0 as i64, |acc, index| { let idx = *index as usize; acc + self.get(idx).unwrap() // 单累加器串行依赖 + 重复 downcast })优化后(核心 inline asm):
core::arch::asm!( "mov z2.d, #0", // accumulator = 0 "ptrue p1.d", // all-true for reduce "whilelt p0.d, {pos}, {n}", // loop predicate "2:", "ld1d {{z0.d}}, p0/z, [{idx}, {pos}, lsl #3]", // load indices "ld1d {{z1.d}}, p0/z, [{val}, z0.d, lsl #3]", // SVE gather (920B-specific) "add z2.d, p0/m, z2.d, z1.d", // predicated accumulate "incd {pos}", // VL-agnostic step "whilelt p0.d, {pos}, {n}", "b.first 2b", // canonical SVE loop tail "uaddv d3, p1, z2.d", // horizontal reduce "fmov {sum}, d3", ... );whilelt+incd+b.first是 SVE 的标准 vector-length-agnostic 循环模式,自动处理尾部不完整向量。u64 内核复用 i64 内核(整数 add 对有无符号一致)。优化二:SVE gather + 4 路独立累加器 + faddv(f64)
f64 求和最大瓶颈是
fadd的 4-5 周期 latency × 单累加器串行依赖。SVE 内核用 4 个独立累加器(z3 / z4 / z5 / z6)打破串行,每 4 个 SVE 向量为一个 chunk 主循环,残余元素用单累加器 tail 循环处理:// Main loop: 4-way unrolled "ld1d {{z0.d}}, p1/z, [{idx}, {pos}, lsl #3]", "ld1d {{z1.d}}, p1/z, [{val}, z0.d, lsl #3]", "fadd z3.d, p1/m, z3.d, z1.d", // chain 0 "incd {pos}", // ... chunks 1, 2, 3 feeding z4, z5, z6 // Tail loop with whilelt predicate // Tree reduce 4 accumulators -> 1 "fadd z3.d, p1/m, z3.d, z4.d", "fadd z5.d, p1/m, z5.d, z6.d", "fadd z3.d, p1/m, z3.d, z5.d", "faddv d7, p1, z3.d", // horizontal reduce4 路独立累加器让 CPU 每个周期可以同时 dispatch 4 个
fadd。接受faddvrecursive pairwise reduction 顺序与 strict left-to-right scalar 求和存在 ULP 级数值差异。优化三:标量 4 路累加器 + PRFM 软件预取(f32 / i128)
f32 和 i128 没有高效的 SVE 路径:
ld1w配合 .d 谓词的 gather 需要复杂的uzp1/fcvt数据重排,throughput 增益不抵指令开销退化为标量 4 路独立累加器 +
PRFM软件预取:while i + 4 <= n { // PRFM with distance 32 for i128 (16 bytes per element) if i + 32 < n { let pf_idx = *group.get_unchecked(i + 32) as usize; prfm_l1_strm(values.add(pf_idx)); // pldl1strm: streaming hint } s0 = s0.wrapping_add(*values.add(group[i ] as usize)); s1 = s1.wrapping_add(*values.add(group[i+1] as usize)); s2 = s2.wrapping_add(*values.add(group[i+2] as usize)); s3 = s3.wrapping_add(*values.add(group[i+3] as usize)); i += 4; }LLVM 会自动把 i128 的
wrapping_add编译成adds + adc配对加法,并使用ldp x, x, [ptr]ARM 配对加载指令。优化四:跨 group
pldl1keep预取在
grouped_sum的外层 group 循环中,对未来第 4 个 group 的 indices 头部地址做软件预取:for (gi, g) in groups.iter().enumerate() { if gi + 4 < n_groups { unsafe { sve_sum::prefetch_next_group(&groups[gi + 4]); // prfm pldl1keep } } // process current group via SVE kernel }pldl1keep(temporal hint)用于 indices 数据,因为它会在内层 SVE 循环中被重复读取,应保留在 L1。value 数据数组的预取在 SVE 内核内部由pldl1strm(streaming hint)处理。优化五:消除
as_arrow()重复 downcast +Vec::with_capacity预分配把原
groups.iter().map(|g| g.iter().fold(...))模式改写为显式两层循环,外层一次性as_arrow()+Vec::with_capacity预分配:优化前:
DataArray::<$T>::from_field_and_values( self.field.clone(), groups.iter().map(|g| { g.iter().fold(0, |acc, idx| acc + self.get(*idx as usize).unwrap()) // ^^^ 每次都做 Any::downcast_ref }), )优化后:
let arrow_array = self.as_arrow()?; let values_ptr = arrow_array.values().as_ptr(); let nulls = if self.null_count() > 0 { arrow_array.nulls() } else { None }; let mut out: Vec<Option<$AggType>> = Vec::with_capacity(n_groups); for (gi, g) in groups.iter().enumerate() { let v = unsafe { sve_kernel(values_ptr, g) }; out.push(Some(v)); } DataArray::<$T>::from_iter(self.field.clone(), out.into_iter())优化六:null 路径 4 路标量累加器 + any 标志
原始 null 路径用
Option<AggType>+ 4-way match 累加,每次迭代都要做Option模式匹配 + 两个分支预测。改写为独立的 4 路累加器 +any: bool标志:let mut any = false; let (mut a0, mut a1, mut a2, mut a3) = (0, 0, 0, 0); while i + 4 <= n { if nulls.is_valid(i0) { a0 = a0.wrapping_add(*values.add(i0)); any = true; } if nulls.is_valid(i1) { a1 = a1.wrapping_add(*values.add(i1)); any = true; } if nulls.is_valid(i2) { a2 = a2.wrapping_add(*values.add(i2)); any = true; } if nulls.is_valid(i3) { a3 = a3.wrapping_add(*values.add(i3)); any = true; } i += 4; } if any { Some(a0.wrapping_add(a1).wrapping_add(a2).wrapping_add(a3)) } else { None }is_valid(idx)编译为 bit-test,分支预测器对 "non-null" 占主导的 TPC-H 数据高度准确。优化推广
所有 6 项优化通过
impl_daft_numeric_agg!宏统一应用到 5 个数值类型(Int64Type/UInt64Type/Float32Type/Float64Type/Decimal128Type)。aarch64 路径用#[cfg(target_arch = "aarch64")]与原 fold 实现隔离,x86 上完全保留原实现,零回归风险。3.2 Story业务交互流程(必填)
无具体业务交互。
3.3 接口设计
外部调用接口保持一致,
DaftSumAggable::grouped_sum公开 trait 签名不变。新增内部内核:
sve_grouped_sum_i64/sve_grouped_sum_i64_with_nullssve_grouped_sum_u64/sve_grouped_sum_u64_with_nullssve_grouped_sum_f64/sve_grouped_sum_f64_with_nullssve_grouped_sum_f32/sve_grouped_sum_f32_with_nullssve_grouped_sum_i128/sve_grouped_sum_i128_with_nullsprfm_l1_keep/prfm_l1_strm/prfm_l2_keep/prfm_l2_strm软件预取助手所有新增内核位于
src/daft-core/src/array/ops/aarch64/目录,由#![cfg(target_arch = "aarch64")]全局门控。3.4 DFX可定位性设计
无
3.5 升级兼容性
无
3.6 性能分析
在 TPC-H 基准测试上,针对 Daft v0.7.5 进行编译测试,得到下面的测试结果:
性能提升:跨平台性能比从 0.818838922 提升至 0.825198708,提升约 0.78%。
3.7 SFMEA分析、测试设计
3.7.1 测试设计
本次优化复用现有测试用例进行验证,无需编写新测试,另外针对 SVE 内核新增 Rust 单元测试。
现有测试覆盖:
tests/dataframe/test_aggregations.pytests/dataframe/test_decimals.pydf.sum()全表求和 /groupby(...).sum()分组求和tests/dataframe/test_morsels.pydf.groupby(...).sum()多分区分组求和tests/dataframe/test_pivot.pytests/dataframe/test_monotonically_increasing_id.py新增单元测试:
src/daft-core/src/array/ops/aarch64/tests.rs测试命令:
安装优化后 wheel 进行以下测试:
# DataFrame 级别聚合测试(覆盖 grouped_sum 的端到端调用路径) pytest tests/dataframe -k "not lance" # Rust 层级 SVE 内核单元测试(aarch64 主机上执行) cargo test -p daft-core --release -- aarch64::tests关键验证点:
test_aggregations.py测试通过 → 验证 i64 / u64 / f64 / f32 各类型 grouped_sum 正确性test_decimals.py测试通过 → 验证 i128 (Decimal128) grouped_sum 正确性test_morsels.py测试通过 → 验证多分区下 grouped_sum 行为一致Describe alternatives you've considered
Additional Context
Would you like to implement a fix?
Yes