用于分析的纯 Clojure 列式数据库
摘要
Flatiron 是一个纯 Clojure 的列式分析库,用于内存表,具有类似 SQL 的 DSL,专为使用原始数组和批处理的高性能而设计。
查看缓存全文
缓存时间: 2026/06/12 18:56
yogthos/flatiron 源代码:https://github.com/yogthos/flatiron
Flatiron
Flatiron 是一个用于 Clojure 的列式分析库。它允许你使用类似 SQL 的 DSL 在内存表上快速运行分析查询,并能对同一数据执行图算法。它是纯 Clojure 实现的,除了 core.async 之外没有其他依赖。可以把它看作是在不想引入完整嵌入式数据库时的选择:加载一些数据,运行分组聚合、排序和过滤,甚至可能在图中运行 PageRank,所有这些都在进程内完成,零配置。
为什么选择列式存储
大多数 Clojure 程序将表格数据表示为 map 的序列。这对于几千行数据来说没问题,但在更大的数据集上就会失效:每一行都是一个堆分配的 map,每个值都是装箱的,每次访问都要通过多层间接。Flatiron 将数据存储为类型化的原始数组——每列一个数组。整数列是 long[],浮点数列是 double[],以此类推。操作直接在这些数组上使用未经检查的算术进行循环,JVM 可以将其优化为紧凑的原生代码。空值使用哨兵值处理,而不是装箱类型,因此没有指针追踪。Morsel 引擎以 1024 行为批次处理数据。这摊销了类型分发的开销:每批决定一次操作,然后在原始类型上运行紧密循环。结果是性能接近原生 C,而不是惯用的 Clojure。
安装
将 git 依赖添加到你的 deps.edn:
{io.github.yogthos/flatiron {:git/tag "v0.2.0" :git/sha "98d700ee79b5425cd837db5b7866a69cf4a0f432"}}
它依赖于 Clojure 1.12.0 和 core.async 1.6.681,并且需要 JDK 18+(哈希内核使用 Math/unsignedMultiplyHigh);CI 在 21 和 25 上运行。
概念
列和表
有五种列类型。每种都将数据存储为 Java 原始数组,并带有可选的空值哨兵:
- I64 — 64 位有符号整数(
long[]) - F64 — 64 位浮点数(
double[]) - Bool — 布尔值(
byte[]) - Sym — Clojure 关键字(
Object[]) - Str — 字符串(
Object[])
Table 是一个模式(关键字列名的向量)加上一个列的向量。仅此而已——没有元数据,没有索引,只有带名称的原始数组。
(require '[flatiron.column :as col])
(require '[flatiron.table :as tbl])
(let [dragons (col/sym-column [:smaug :fafnir :tiamat :smaug :fafnir])
gold (col/i64-column [9000 750 1200 3100 2400])
table (tbl/table [:Dragon :Gold] [dragons gold])]
(tbl/nrows table) ;; => 5
(tbl/ncols table) ;; => 2
(tbl/col table :Gold)) ;; => #<I64Column [9000 750 1200 3100 2400]>
;; — 史矛革又在囤积了
自定义类型
像日期和时间戳这样的领域类型不需要自己的列类型或一个 Object[] 来装箱每个值。一个 LocalDate 就是它的纪元日,一个 Instant 就是它的纪元毫秒,这些编码保留了顺序,因此比较、排序、分组和 min/max 在底层的原始类型上就已经是正确的。Flatiron 将这样的类型存储为普通的 long[] 列,并标记一个逻辑类型:列向其每个操作报告其物理类型(:i64),因此热循环不变,值永远不会被装箱,只有在边界(构建列、读取值、持久化)时才运行编解码器,将领域对象转换进来或转换出去。
(require '[flatiron.column :as col])
(let [hired (col/date-column [(java.time.LocalDate/of 2019 4 1)
(java.time.LocalDate/of 2021 9 15)])]
(col/-type-tag hired) ;; => :i64 (物理——操作分发的依据)
(col/-logical-tag hired) ;; => :date (逻辑——值如何暴露)
(col/-get-obj hired 0)) ;; => #object[java.time.LocalDate "2019-04-01"]
内置的逻辑类型,全部由 long[] 支持::date(LocalDate)、:instant(Instant)、:datetime(LocalDateTime)、:date-millis(java.util.Date)和 :duration(Duration)。使用 col/typed-column 或 date-column/instant-column/datetime-column/duration-column 辅助函数为其中任何一个构建列。
谓词直接接受领域值。字面量在每个谓词中被编码一次,在行循环之外,因此比较仍是一个原始操作:
(-> employees
(where (>= :Hired (java.time.LocalDate/of 2020 1 1)))
(select :Dept (count :Hired)))
保留值的操作会保留逻辑类型:过滤、排序、分组键以及日期列的 min/max 都返回日期。产生新数字的聚合会丢弃它,因此对日期列进行 sum 和 avg 会返回普通的 :i64/:f64。二进制存储在其元数据中记录了逻辑类型,因此保存的表在往返过程中会保留其类型。
使用 flatiron.types/register-type! 注册你自己的类型,为其指定物理支持(:i64 或 :f64)以及编码/解码对:
(require '[flatiron.types :as types])
(types/register-type!
:cents
{:physical :i64
:class java.math.BigDecimal
:encode (fn ^long [^java.math.BigDecimal d] (.longValueExact (.movePointRight d 2)))
:decode (fn [^long v] (.movePointLeft (java.math.BigDecimal/valueOf v) 2))})
(col/typed-column :cents [(bigdec "1.50") (bigdec "2.25")])
过滤
where 宏通过将列与常量比较来构建布尔掩码,使用 and、or 和 not 组合复合谓词的掩码,然后将通过的行物化到一个新表中。一个三级选择位图存在于 flatiron.selection 中,作为更底层的原语,供想要跟踪哪些行存活而不物化的调用者使用,但内置的 where 会急切地物化。对于过滤然后聚合的管道,flatiron.group/group-by 和 parallel-group-by 通过 :where 选项直接接受掩码,并且只通过该掩码收集键和聚合列,完全跳过中间表。
Morsel 引擎
以 Rayforce 中的“morsel”(一小口数据)概念命名。逐元素操作(算术、比较)从列创建 morsel 源,并通过它拉取 1024 行的批次;在每个批次内,循环体直接对原始数组运行,没有协议分发。聚合和分组更进一步:它们在类型特化的循环中直接读取列的后备数组,仅在不必要时才回退到 morsel 层。无论哪种方式,你都得到类型化泛型操作的抽象,同时获得手写原始循环的性能。
DSL
DSL 在宏展开时编译为基于 morsel 的操作。没有运行时查询解析——宏直接发出函数调用。
(require '[flatiron.dsl :refer [sum count avg min max]])
;; 分组聚合
(select trades :Symbol (sum :Qty))
;; 多个聚合函数
(select trades :Symbol (sum :Qty) (avg :Price) (count :Qty))
;; 多个分组键
(select trades :Region :Side (sum :Qty))
;; 交叉表
(pivot trades :Symbol :Side :Qty sum)
过滤器
(-> trades
(where (> :Qty 100))
(select :Symbol (sum :Qty)))
这将筛选出 Qty > 100 的行,然后按 Symbol 分组并对 Qty 求和。where 宏将谓词编译为对列运行时类型分发的掩码构建调用,因此与浮点列比较的整数字面量会自动进行强制转换,而不是从字面量的类型中选择比较函数。支持的谓词:>、<、>=、<=、=、not=,以及通过 and、or 和 not 组合。
聚合函数
所有聚合都是单次遍历的归约,它们对列类型进行一次分发,然后直接对后备原始数组使用未经检查的算术进行循环。
| 函数 | 描述 |
|---|---|
sum | 值的和,跳过空值 |
count | 非空值的计数 |
avg | 算术平均值,跳过空值;对于没有非空值的组返回 null |
min | 最小值,跳过空值;对于没有非空值的组返回 null |
max | 最大值,跳过空值;对于没有非空值的组返回 null |
排序和窗口函数
排序在索引数组上使用 java.util.TimSort——列数据从不被重新排列。排序是稳定的,支持升序和降序。
(require '[flatiron.sort :as sort])
(require '[flatiron.window :as win])
(let [sorted (sort/sort-table table [[:Qty :asc]])]
(win/row-number sorted)) ;; => I64Column [1 2 3 4 5]
窗口函数在排序列上操作:
row-number— 从 1 开始的连续编号rank— 带间隙的排名(1, 1, 3, 3, 5)dense-rank— 无间隙的排名(1, 1, 2, 2, 3)lag/lead— 访问上一行或下一行的值,带有偏移量和默认值
并行执行
主要的并行入口点是 flatiron.group/parallel-group-by:它通过键哈希的高位对行进行基数分区,然后在共享的 ForkJoin 池上运行哈希、直方图、分散和每个分区的分组阶段。分区是不相交的,因此结果无需合并步骤即可连接,输出与单线程的 group-by 相同。
(require '[flatiron.group :as g])
(g/parallel-group-by table
:keys [:Region]
:aggs [{:agg :sum :col :Qty :out :total}]
:n-threads 8)
flatiron.parallel 还提供了基于 core.async 线程的每列并行原语(parallel-i64-sum、parallel-i64-min、并行过滤计数等)。当每行有足够的工作时,并行化才有回报——普通的标量求和受限于内存带宽,单线程运行得一样快。
图算法
Flatiron 包含一个 CSR(压缩稀疏行)图引擎。你可以从两列(源节点 ID 和目标节点 ID)构建图,它会在单次遍历中构造前向和反向邻接结构。
(require '[flatiron.graph :as g])
(let [src (col/i64-column [0 0 1 2 3])
dst (col/i64-column [1 2 3 3 0])
graph (g/graph src dst)
result (g/page-rank graph 20 0.85)]
;; result 是一个包含 :node 和 :rank 列的表
)
可用算法:
- BFS — 从起始节点的广度优先搜索
- DFS — 深度优先搜索,迭代(非递归)
- Dijkstra — 带权重的单源最短路径
- PageRank — 可配置阻尼因子和迭代次数的迭代 PageRank
- Connected components — 通过 BFS 查找弱连通分量
I/O
CSV
CSV 读取器通过采样前 100 行进行类型推断,然后读取到类型化列中。它能自动处理广泛的类型。read-csv 接受 CSV 内容作为字符串或 Reader(而不是文件路径):
(require '[flatiron.io :as io])
(require '[clojure.java.io :as jio])
(with-open [r (jio/reader "data/trades.csv")]
(let [table (io/read-csv r)]
(select table :Symbol (sum :Qty))))
写入器将表输出回 CSV。
二进制列式存储
为了实现快速持久化,Flatiron 有一种每列一个文件的二进制格式。每个列变成单个文件(col_Symbol、col_Qty、…),并附带一个小型的 _meta.edn 描述模式。读取时尽可能实现零拷贝。
(require '[flatiron.store :as store])
(store/save-table table "data/trades_store")
(let [loaded (store/load-table "data/trades_store")]
(select loaded :Symbol (sum :Qty)))
为什么选择 Flatiron 而非其他替代方案
有几个优秀的 Clojure 数据库。Flatiron 填补了一个特定的细分市场:
- 与
clojure.core序列操作比较 — 在大量行的数值聚合上,Flatiron 快两到三个数量级,但它只适用于自己的列类型。快速脚本使用核心函数,繁重的工作使用 Flatiron。 - 与
tech.ml.dataset比较 — TMD 功能更丰富(日期处理、统计函数、与多种格式的互操作、完整的 I/O 生态系统)。Flatiron 更小——它除了 Clojure 本身之外唯一的依赖是 core.async——并专注于更窄操作集上的原始速度。 - 与嵌入式数据库(H2、SQLite)比较 — 数据库提供 SQL、事务和持久化。Flatiron 提供进程内数据,你可以直接从 Clojure 操作,无需通过 JDBC。如果你已经在 Clojure 数据结构中有数据,只需要快速分析,Flatiron 的仪式感更少。它也是一个图引擎:从两列构建 CSR 图,并在你正在聚合的相同数据上运行 BFS、Dijkstra、PageRank 或连通分量——在 SQL 数据库中,这最多意味着递归 CTE,或导出到单独的图库。
与 Rayforce(C)的基准测试
bench/flatiron/rayforce_bench.clj 通过 Flatiron 和通过 Rayforce(https://github.com/RayforceDB/rayforce)(Flatiron 重新实现的 C17 引擎)运行相同的查询。两个引擎读取相同生成的数据集;Rayforce 使用其内置的 timeit 计时,Flatiron 使用 System/nanoTime,取预热后至少 10 次运行的最小值。
clojure -M:bench -m flatiron.rayforce-bench [path-to-rayforce-binary]
结果在 Apple M1 Max(JDK 26,Rayforce make release),1M 行。flatiron 列是单线程路径;flatiron par8 是使用 :n-threads 8 运行的并行路径(对于分组查询使用 parallel-group-by,对标量求和使用 parallel-i64-sum)。Rayforce 使用自己的工作池进行内部并行化。
| 查询 | rayforce (C) | flatiron | flatiron par8 |
|---|---|---|---|
| 按 Sym 分组(100 组),求和 | 0.99 ms | 21.5 ms | 6.7 ms |
| 按 Sym 分组,sum+count+avg | 1.34 ms | 23.9 ms | 11.5 ms |
| 过滤 Qty>500,按 Sym 分组求和 | 0.45 ms | 14.8 ms | 7.2 ms |
| 按 Id 分组(100K 组),求和 | 2.83 ms | 33.2 ms | 8.4 ms |
| 标量求和,1M i64 | 0.05 ms | 0.7 ms | 1.0 ms |
C 实现比 Flatiron 的并行路径快 3–10 倍(它使用 SIMD 内核、自定义分配器,并在扫描时饱和内存带宽)。Flatiron 的目标是在保持纯 Clojure 的同时,性能保持在 C 的一个数量级以内。标量求和是一个警示行:它受限于内存带宽,因此并行版本输给了单线程循环——线程分发的开销超过了它节省的时间。
致谢
Flatiron 是来自 Rayforce(https://github.com/RayforceDB/rayforce)思想的 Clojure 重新实现,Rayforce 是一个用 C17 编写的 SIMD 加速列式分析和图引擎。Morsel 驱动的执行模型、1024 元素的批次大小、CSR 图布局及其 BFS、DFS、Dijkstra 和 PageRank 算法,以及类似 SQL 的表面都遵循了 Rayforce 的设计。flatiron.hash 中的 wyhash 哈希移植自 Rayforce 的 src/ops/hash.h,DSL 借用了其 Rayfall 查询语言。Rayforce 是 MIT 许可的。
许可
MIT
相似文章
RayforceDB – 一个具有类似Lisp语法的纯C分析数据库
RayforceDB是一个开源的、纯C编写的分析数据库,它将列式分析、图遍历和递归查询结合成一个单一的嵌入式管道,用于高性能、低延迟的数据处理。
DuckDB – 笔记本电脑上的数据强力工具,现已支持 Clojure(2023)
TechAscent 展示了如何通过 tech.ml.dataset(TMD)从 Clojure 中利用 DuckDB 的高性能矢量化 SQL 引擎,实现大型内存外连接,以及约两分钟内将 50GB CSV 数据摄入并压缩至 18GB。
Fluree DB(GitHub 仓库)
Fluree DB 是一个开源的时间图数据库,具有类似 Git 的分支、集成的向量/文本/地理搜索、细粒度的访问控制,并支持 SPARQL、JSON-LD 和 Open Cypher。它针对 AI 代理记忆进行了优化,在十亿级图上实现了高性能。
列式存储即规范化
本文将列式存储重新定义为数据库规范化的极端形式,展示了把属性拆分为位置对齐的数组如何与基于隐式序数主键连接的规范化表如出一辙。
Biff.graph:将你的 Clojure 代码库构建为可查询图
Biff 是一个面向独立开发者的开源 Clojure Web 框架,提供构建 Web 应用的工具。