Mini-Arrow:极简的Rust实现的Arrow

本文介绍 mini-arrow 项目,它是一个Rust编写的迷你Apache Arrow项目的实现,仅用1600行的Rust代码,展示了Arrow中几个最核心问题的设计和解决方案。项目参考了Type Exercise in Rust

在展开讲解之前,我们先了解Arrow项目:

Arrow 要解决什么问题 #

Arrow 是一个列式内存格式项目,它要解决以下几个核心问题:

行式 vs 列式 #

传统数据库(如 MySQL)按存储数据:一条记录的所有字段连续放在一起。

行式存储:
  [id=1, name="Alice", age=30] [id=2, name="Bob", age=25] [id=3, ...]

而列式存储(如 Arrow、ClickHouse)按存储:同一列的所有值连续放在一起。

列式存储:
  id:    [1, 2, 3, ...]
  name:  ["Alice", "Bob", ...]
  age:   [30, 25, ...]

为什么列式更好? 分析型查询(如 SELECT AVG(age) FROM users)只关心 age 一列。列式存储能:

  1. 只读需要的列,跳过无关数据,减少 I/O;
  2. 缓存局部性好——同一列的值在内存里连续,扫描时命中率高;
  3. 支持向量化——对连续内存做 SIMD 批量运算,比逐行解释快几个数量级。

列式存储的难点 #

把数据按列存放,会立刻遇到几个棘手问题:

难点一:定长 vs 变长。 i32 可以放进 Vec<i32>,但 String 长度不定,不能直接放进一个定长数组。需要额外的布局来管理变长数据。

难点二:空值(NULL)。 Vec<T> 要求连续、定长,无法"跳过"某个位置。NULL 怎么表示?如果为每个值包一层 Option,会破坏连续布局、浪费内存、无法向量化。

难点三:owned 值 vs 零拷贝引用。 读取 StringArray 时,我们希望直接拿到 &str 而不拷贝;但构建数组时又需要持有 String 的所有权。两种形态如何共存?

难点四:运行时才知道类型。 Rust 的泛型在编译期展开,但数据库引擎在运行时才知道列的类型(从 SQL 解析出来)。如何让任意类型的数组被统一传递、分发、求值?

难点五:表达式求值的重复决策。 每个二元运算(+<=contains)都要处理参数个数、类型检查、长度检查、NULL 传播、输出构建。如果这些逻辑写死在每个循环里,每新增一个函数就多一处重复和出错的机会。

难点六:类型族扩展。 每加一个类型(如 f64),就要重复物理变体、标量变体、数组别名、构建器、转换……重复越多,越容易漂移。

mini-arrow 的四个模块——scalararraybuilderexpr——正是为逐一攻克这些难点而设计的。下面从"一个看似简单的问题"开始,逐步拆解每个模块的设计动机。

起点:一个看似简单的问题 #

手写一个 i32 + i32 的循环很容易:

for row in 0..left.len() {
    output.push(match (left.get(row), right.get(row)) {
        (Some(left), Some(right)) => Some(left.wrapping_add(right)),
        _ => None,
    });
}

这段代码对 i32 没问题。但一旦要求更多,它就开始“漏气”:

  1. 字符串怎么办? StringArray::get 如果返回 String,每次都要拷贝;我们希望返回 &str
  2. NULL 怎么办? 不能返回 Vec<Option<T>>,那会破坏列式布局。
  3. 混合类型怎么办? i32 <= i64 该先提升成 i64 再比较。
  4. 运行时才知道类型怎么办? SQL 解析出的函数名是字符串,运行时才能映射到具体实现。
  5. 每加一个函数都要重写一次 arity / 长度 / NULL 检查怎么办?

mini-arrow 的解法是把这些问题拆开,每个模块只解决一个问题:把那些决策从行循环里移出来,沉淀成 scalararraybuilderexpr 四个模块。每个模块都对应开头列出的一个难点,下面逐一拆解:

模块拆解:每个模块在解决什么问题 #

1、 types:划定“基本类型”的边界 #

pub trait PrimitiveType: Default + Copy + Debug + 'static {}
impl PrimitiveType for i32 {}
impl PrimitiveType for i64 {}
impl PrimitiveType for bool {}

要解决的问题:哪些类型可以用“定长数组”存储?

PrimitiveType 是一个标记 trait。满足它的类型才能放进 PrimitiveArray<T>——因为它们 Default + Copy + Debug + 'staticString 不是 Copy,所以被排除在外,必须走 StringArray 这套变长布局。

graph LR
    subgraph PrimitiveType 边界
        PT[PrimitiveType trait] --> |i32, i64, bool| PA[PrimitiveArray&lt;T&gt;]
        PT -.->|String 不满足 Copy| SA[StringArray 变长布局]
    end

2、scalar:owned 值与零拷贝引用 #

要解决的问题 #

数据库里一个"值"有两种形态:

  • 拥有所有权的值(owned):如 String,可以自由持有、移动。
  • 零拷贝引用(borrowed):如 &str,只是借用数组里的字节,不产生分配。

如果强制统一成一种形态,就会出问题:

  • 若统一用 owned,那么 StringArray::get 每次都要 to_string() 拷贝一份,性能灾难。
  • 若统一用引用,那么构建数组时无法持有数据。

为什么这样设计 #

两个 trait 表达这两种形态,并用**泛型关联类型(GAT)**把它们关联起来:

pub trait Scalar: 'static + Clone + Debug + TryFrom<ScalarImpl> + Into<ScalarImpl> {
  type ArrayType: Array<Item = Self>;
  type RefType<'a>: ScalarRef<'a, ScalarType = Self, ArrayType = Self::ArrayType>;
  fn as_scalar_ref(&self) -> Self::RefType<'_>;
}

pub trait ScalarRef<'a>:
  'a + Clone + Copy + Debug + TryFrom<ScalarRefImpl<'a>> + Into<ScalarRefImpl<'a>>
{
  type ArrayType: Array<RefItem<'a> = Self>;
  type ScalarType: Scalar<RefType<'a> = Self>;
  fn as_scalar(&self) -> Self::ScalarType;
}

这两个 trait 形成一组互反箭头

Scalar ──RefType<'a>──> ScalarRef<'a>
  │                         │
ArrayType                ArrayType
  ▼                         ▼
Array ─────Builder─────> ArrayBuilder

关键洞察:基本类型既是 Scalar 又是 ScalarRef

  • i32Copy 的,拷贝它没有成本,所以 RefType<'a> = i32as_scalar_ref 就是 *self
  • String 则不同:ScalarStringScalarRef&'a stras_scalar_ref 返回 self.as_str()
graph LR
    subgraph 整数 i32
        S1[i32 as Scalar] -->|RefType = i32| SR1[i32 as ScalarRef]
        SR1 -->|ScalarType = i32| S1
    end

    subgraph 字符串 String
        S2[String as Scalar] -->|RefType = &str| SR2[&str as ScalarRef]
        SR2 -->|ScalarType = String| S2
    end

为什么需要生命周期索引的关联类型? 因为 StringArray::get 返回的 &str 的生命周期必须与数组的借用绑定。 普通关联类型无法表达"返回值借用 self",必须用 RefItem<'a> 这种 GAT。 for<'a> 这个 HRTB 绑定意味着"对调用者选择的任意借用生命周期都成立"——整数实现可以忽略它,字符串实现不能。

3、array:列式存储与空值 #

要解决的问题 #

列式存储的核心是把同一列的数据连续存放,以获得缓存局部性和向量化能力。但有两个难题:

  1. 定长 vs 变长i32 可以放进 Vec<i32>,但 String 长度不定,不能直接放。
  2. 空值(NULL)Vec<T> 要求连续、定长,无法"跳过"某个位置。

为什么这样设计 #

采用 Arrow 标准布局mini-arrow 完整实现了两种:

定长类型 PrimitiveArray<T> —— data + bitmap

pub struct PrimitiveArray<T: PrimitiveType> {
  data: Vec<T>,      // 连续数据
  bitmap: BitVec,    // 空值位图
}
  • NULL 位置用 T::default() 占位,位图第 i 位标记是否有效。
  • 读取:if bitmap[i] { Some(data[i]) } else { None }
graph LR
    subgraph PrimitiveArray_i32
        data["data: 1, 2, 3, 0, 5"]
        bitmap["bitmap: 1, 1, 1, 0, 1"]
    end
    data -->|第 i 行对应 bitmap 第 i 位| bitmap
    data -->|bitmap 第 3 位为 0| null["第 3 行是 NULL, data 中 0 只是占位"]

变长类型 StringArray —— data + offsets + bitmap

pub struct StringArray {
  data: Vec<u8>,       // 所有字符串拼接的字节
  offsets: Vec<usize>, // 每个字符串的起始偏移
  bitmap: BitVec,      // 空值位图
}
  • i 个字符串的区间是 offsets[i]..offsets[i+1],长度 = 差值。
  • 读取:&data[offsets[i]..offsets[i+1]]零拷贝返回 &str
  • offsets 数的是 UTF-8 字节,不是字符;空行和 NULL 行都重复前一个 offset。
graph LR
    subgraph StringArray
        data["data bytes: h e l l o w o r l d"]
        offsets["offsets: 0, 5, 10"]
        bitmap["bitmap: 1, 1"]
    end
    data -->|所有字节连续存放| offsets
    offsets -->|第 0 行范围 0 到 5| str0["hello"]
    offsets -->|第 1 行范围 5 到 10| str1["world"]
    bitmap -->|两行均有效| offsets

为什么 NULL 用位图而不是 Option<T> 可空性是值状态,不是类型变体。用 Option 或 validity 位图表达,而不是为每个值包一层 Option。 因为 Vec<Option<T>> 会破坏连续布局、浪费内存,且无法向量化。

统一抽象 Array trait

pub trait Array: Sized + 'static + TryFrom<ArrayImpl> + Into<ArrayImpl> {
  type Builder: ArrayBuilder<Array = Self>;
  type Item: Scalar<ArrayType = Self>;
  type RefItem<'a>: ScalarRef<'a, ScalarType = Self::Item, ArrayType = Self>;
  fn get(&self, idx: usize) -> Option<Self::RefItem<'_>>;
  fn len(&self) -> usize;
  fn iter(&self) -> ArrayIterator<Self>;
  fn from_slice(data: &[Option<Self::RefItem<'_>>]) -> Self; // 默认实现:走 builder
}
数组 Item(owned) RefItem(引用)
I32Array i32 i32(Copy,无成本)
StringArray String &str(零拷贝)

这样,一个泛型契约就能同时读取普通、空、全 NULL 的数组,且字符串读取不产生分配。

4、builder:不可变数组的增量构建 #

要解决的问题 #

数组一旦建好就是不可变的(列式存储通常如此)。但构建过程需要增量 push,最后一次性产出。

为什么这样设计 #

把构建抽象成 ArrayBuilder trait,与 Array 通过关联类型互相关联:

pub trait ArrayBuilder {
  type Array: Array<Builder = Self>;
  fn with_capacity(capacity: usize) -> Self;
  fn push(&mut self, value: Option<<Self::Array as Array>::RefItem<'_>>);
  fn finish(self) -> Self::Array;
}

PrimitiveArrayBuilder<T>Some(v) 数据入 data、位图置 trueNone 推入 T::default() 占位、位图置 false

StringArrayBuilderSome(s) 字节追加到 data、记录新 offset;None offset 不变(复用前一个)、位图置 false

graph LR
    subgraph Builder 流程
        WC[with_capacity] --> P1[push Some]
        P1 --> P2[push None]
        P2 --> P3[push Some]
        P3 --> F[finish]
        F --> A[Array 不可变]
    end

为什么 builder 和 array 是"互反"的? Array::Builder 指向能构建它的 builder,ArrayBuilder::Array 指向它产出的数组。 Array::BuilderArrayBuilder::Array 这两个关联类型方向相反,形成“互反”关系。这样 from_slice 的默认实现才能泛型地写出:

let mut builder = Self::Builder::with_capacity(data.len());
// ... push
builder.finish()

5、类型擦除:跨越运行时边界 #

要解决的问题 #

Rust 的泛型在编译期展开,但数据库引擎在运行时才知道列的类型(比如从 SQL 解析出来)。 我们需要一个统一类型,让任意数组都能被传递、分发、求值。

为什么这样设计 #

用三个擦除枚举做类型擦除:

pub enum ArrayImpl {
  Int32(I32Array), Int64(I64Array), Bool(BoolArray), String(StringArray),
}
pub enum ScalarImpl { Int32(i32), Int64(i64), Bool(bool), String(String) }
pub enum ScalarRefImpl<'a> { Int32(i32), Int64(i64), Bool(bool), String(&'a str) }

装箱(upcast)用 From,拆箱(downcast)用 TryFrom

impl From<I32Array> for ArrayImpl {          // 装箱:具体 → 枚举
  fn from(array: I32Array) -> Self { ArrayImpl::Int32(array) }
}

impl TryFrom<ArrayImpl> for I32Array {       // 拆箱:枚举 → 具体,类型不符报错
  type Error = TypeMismatch;
  fn try_from(array: ArrayImpl) -> Result<Self, Self::Error> {
    match array {
      ArrayImpl::Int32(array) => Ok(array),
      other => Err(TypeMismatch(stringify!(Int32), other.identifier())),
    }
  }
}
graph LR
    subgraph 类型擦除
        A1[I32Array] -->|From| E1[ArrayImpl::Int32]
        S1[i32] -->|From| E2[ScalarImpl::Int32]
        R1[&str] -->|From| E3[ScalarRefImpl::String]
    end
    E1 -->|TryFrom| A1
    E1 -.->|类型不符| Err[Err&lt;TypeMismatch&gt;]

为什么拆箱是"可失败的"? 因为 ArrayImpl 是运行时值,它可能装着 StringArray,而调用方请求的是 I32Array。 编译期无法保证这一点,所以必须用 TryFrom 在边界处做受检转换。 这正是 TypeMismatch 错误存在的意义——它把"类型不符"变成可处理的错误,而不是崩溃。 错误的变体返回 Err,而不是 panic。

6、expr:把决策移出行循环 #

要解决的问题 #

每个二元/一元循环都必须遵守的五条重复决策

  1. 检查输入个数(arity);
  2. 检查物理类型;
  3. 检查每个输入长度;
  4. 任一严格输入为 NULL 时跳过标量函数
  5. 构建关联的输出数组,或返回第一个行错误而不产生部分输出。

如果这些决策写死在每个循环里,每新增一个函数就多一处重复和出错的机会。

为什么这样设计 #

泛型适配器把"行循环"和"具体标量函数"分离:

pub trait BinaryExprFunc<I1: Array, I2: Array, O: Array> {
  fn eval<'a>(&self, i1: I1::RefItem<'a>, i2: I2::RefItem<'a>) -> O::Item;
}

BinaryExpression<I1, I2, O, F> 是泛型二元表达式,eval_batch 只写一次行循环:

pub fn eval_batch(&self, i1: &ArrayImpl, i2: &ArrayImpl) -> Result<ArrayImpl> {
  let i1a: &I1 = i1.try_into()?;   // 类型擦除 → 具体类型
  let i2a: &I2 = i2.try_into()?;
  assert_eq!(i1.len(), i2.len(), "array length mismatch");
  let mut builder: O::Builder = O::Builder::with_capacity(i1.len());
  for (i1, i2) in i1a.iter().zip(i2a.iter()) {
    match (i1, i2) {
      (Some(i1), Some(i2)) => builder.push(Some(self.expr.eval(i1, i2).as_scalar_ref())),
      _ => builder.push(None),   // NULL 传播
    }
  }
  Ok(builder.finish().into())
}
graph LR
    subgraph 表达式求值流程
        A1[ArrayImpl] -->|try_into| I1[&I1]
        A2[ArrayImpl] -->|try_into| I2[&I2]
        I1 --> ZIP[iter zip]
        I2 --> ZIP
        ZIP -->|Some, Some| EVAL[F::eval]
        ZIP -->|含 None| NULL[push None]
        EVAL --> PUSH[builder.push]
        NULL --> PUSH
        PUSH --> FINISH[finish]
        FINISH --> OUT[ArrayImpl]
    end

NULL 传播是 SQL 三值逻辑的核心:任一输入为 NULL,输出就是 NULL,且不调用标量函数

具体函数通过 BinaryExprFunc 实现,例如比较运算支持类型提升

pub struct ExprCmpLe<I1: Array, I2: Array, C: Array>(pub PhantomData<(I1, I2, C)>);

impl<I1: Array, I2: Array, C: Array> BinaryExprFunc<I1, I2, BoolArray> for ExprCmpLe<I2, I2, C>
where
  for<'a> I1::RefItem<'a>: Into<C::RefItem<'a>>,
  for<'a> I2::RefItem<'a>: Into<C::RefItem<'a>>,
  for<'a> C::RefItem<'a>: PartialOrd,
{
  fn eval<'a>(&self, i1: I1::RefItem<'a>, i2: I2::RefItem<'a>) -> bool {
    i1.into().partial_cmp(&i2.into()).unwrap() == Ordering::Less
  }
}

第三个类型参数 C公共提升类型——ExprCmpLe::<_, _, I64Array> 表示把两个 i32 提升为 i64 再比较。

为什么用 PhantomData ExprCmpLe 本身不存任何数据,I1/I2/C 只用于类型层面。 PhantomData<(I1, I2, C)> 让编译器知道这个类型"拥有"这三个类型参数, 从而正确推导 trait 实现和生命周期,而不占用任何运行时内存。

最后,build_binary_expression运行时工厂,把 ExpressionFunc 枚举映射到具体的 Box<dyn Expression>

pub fn build_binary_expression(f: ExpressionFunc) -> Box<dyn Expression> {
  match f {
    CmpLe => Box::new(BinaryExpression::<I32Array, I32Array, BoolArray, _>::new(
      ExprCmpLe::<_, _, I32Array>(PhantomData),
    )),
    StrContains => Box::new(BinaryExpression::<StringArray, StringArray, BoolArray, _>::new(ExprStrContains)),
  }
}

这解决的是从运行时名字选择具体类型化表达式的问题——SQL 解析器拿到的是字符串函数名,需要映射到具体的类型化实现。

7、宏系统:让类型族可扩展 #

要解决的问题 #

每加一个类型(如 f64),就要重复物理变体、标量变体、数组别名、构建器、转换…… 重复越多,漂移越可能发生——可能加了数组却忘了加转换,或转换里写错类型名。

为什么这样设计 #

一个类型族目录(catalog),让"一行定义一个类型族"。 mini-arrowfor_all_variants! 宏实现:

macro_rules! for_all_variants {
  ($macro:ident $(, $x:ident)*) => {
    $macro! {
      [$($x),*],
      { Int32, int32, I32Array, I32ArrayBuilder, i32, i32 },
      { String, string, StringArray, StringArrayBuilder, String, &'a str }
    }
  };
}

每个元组是 { 枚举变体名, 函数后缀, 数组类型, 构建器类型, owned 类型, 引用类型 }。 然后 impls.rs 用这个清单批量生成所有样板代码:

  • impl_array_dispatchArrayImplget / len / is_empty / identifier
  • impl_array_conversion:所有 From / TryFrom 转换
  • impl_scalar / impl_scalar_conversionScalar / ScalarRef 实现与转换

为什么用宏而不是手写? 宏让目录成为唯一事实来源——新增类型只需在清单加一行,其余代码自动生成, 且不会出现"加了数组却忘了加转换"这种漂移。省略或重复一个类型族会变成编译或测试失败。

注意mini-arrowmacros.rsimpls.rs 目前被注释掉了 (//mod macros;//mod impls;),实际生效的是手写的 array_impl.rsbuilder_impl.rs 等。 这说明项目正处于从手写代码迁移到宏生成的过渡阶段——宏版本已写好但未启用。

总结 #

模块 要解决的问题 为什么这样设计
scalar owned 值与零拷贝引用如何共存 Scalar / ScalarRef 两个 trait + GAT 关联;基本类型双实现,String/&str 分离
array 定长/变长、空值如何列式存储 Arrow 标准布局:data+bitmap(定长)、data+offsets+bitmap(变长);NULL 是值状态不是类型
builder 不可变数组如何增量构建 ArrayBuilder trait 与 Array 互反关联,from_slice 默认实现走 builder
类型擦除 运行时如何统一分发任意类型 三个擦除枚举 + From/TryFrom 受检转换;错误变体返回 Err 而非 panic
expr 如何避免每个循环重复决策 泛型适配器 BinaryExpression 把行循环与标量函数分离;NULL 传播、类型提升
宏系统 如何让类型族可扩展 for_all_variants 目录一行定义一个类型族,批量生成样板代码

贯穿始终的设计原则:

  1. 把决策移出行循环:arity、类型、长度、NULL、输出构建,都由适配器统一处理。
  2. 可空性是值状态:用 Option / validity 位图,而不是 DataType::Nullable 变体。
  3. 受检的运行时擦除:downcast 是可失败的,错误变体返回 Err
  4. 目录驱动扩展:一行定义一个类型族,省略或重复变成编译/测试失败。

回到开头:难点如何被逐一攻克 #

开头列出的难点 由谁解决
难点一:定长 vs 变长 array 的两种布局:PrimitiveArray(定长)、StringArray(offset + 扁平 buffer)
难点二:空值(NULL) array 的 validity 位图,NULL 是值状态而非类型变体
难点三:owned 值 vs 零拷贝引用 scalarScalar / ScalarRef 双 trait + GAT
难点四:运行时才知道类型 类型擦除枚举 + From / TryFrom 受检转换
难点五:表达式求值的重复决策 expr 的泛型适配器 BinaryExpression
难点六:类型族扩展 宏系统 for_all_variants 目录