普段Parserを作るとき、自分はその構造自体をあまり疑ったことがありませんでした。
Input
↓
Parser
↓
Record
↓
Processing
入力を読み、parseしてrecordを返し、それを後段で処理する。
streamingやbuffering、並列化は考えても基本的にはこの形です。
ところが最近RedditのRustコミュニティでこの前提をひっくり返す面白いCSV Parserが紹介されていました。
csveeeというRust製の並列CSV Parserです。
投稿者の Simon Ellmann 氏は、Thomas Neumann 氏との共著論文の著者本人でもあります。
論文はこちらです。
One Pass to Parse Them All: Fused Parallel CSV Processing
論文では最大 180 GB/s のthroughputを報告しています。
CSVを並列に読むのは意外と面倒
たとえば巨大なCSVなら分割して複数threadで読めばよさそうに見えます。
Chunk A → Thread A
Chunk B → Thread B
Chunk C → Thread C
しかしCSVではquoted fieldの中に改行を含められます。
id,name,address
1,Alice,"Tokyo
Shinjuku"
2,Bob,Sapporo
そのため、ファイル途中から読み始めたworkerは、その位置がrecordの途中なのか、新しいrecordの開始なのかを簡単には判断できません。
csveee ではここに speculative parsing を使います。
開始状態を仮定して各chunkを並列にparseし、後からchunk境界の整合性を確認します。
仮定が間違っていた部分だけ、正しい位置からreparseします。
つまり、
境界を全部調べてから並列化するのではなく、先に並列処理して後から整合性を確定する
という発想です。
Parserからrecordを受け取らない
個人的に一番面白かったのはこちらです。
普通なら、
for record in parser.records() {
process(record);
}
のように、Parserからrecordを受け取って処理します。
csveee はこれを逆にして、処理するコードをParser側へ渡します。
READMEでも、
instead of the parser handing records to your code, you hand your code to the parser
と説明されています。
概念的には、
通常
Parser
↓
Record
↓
Processing
から、
csveee
Parser
└─ Parsing + Processing
へ変える感じです。
これによってparseと後段処理を同じparallel executionの中で行えます。
init / acc / merge
そのための基本形が、
init
acc
merge
です。
init でchunkごとのstateを作り、acc でrecordを処理しながらstateを更新し、最後に merge で統合します。
Chunk A → State A
Chunk B → State B
Chunk C → State C
State A
State B
State C
↓
merge
↓
Result
各workerが局所stateだけを持てるため、巨大な中間collectionを一度作ってから処理する必要がありません。
さらにrecordは可能な限りParser buffer上のsliceとして扱われ、不要なcopyも減らされています。
SIMDより面白いのはexecution model
csveee 自体にはSIMD parserやbounded ring bufferなど、かなり低レイヤーな最適化も入っています。
ただ、自分が特に面白いと思ったのはSIMDそのものではなく、
Parserとprocessingの境界を変えたこと です。
Parserを単体で速くするのではなく、
Parsing
+
Processing
+
Parallelism
を一つのexecution modelとして扱っています。
180GB/sはどう見るべきか
論文では128-core AMD EPYC環境で最大 180 GB/s を報告しています。
Reddit上の追加測定では、TPC-H lineitem約80GBに対して次の結果が紹介されています。
| Parser | Throughput |
|---|---|
| csveee SIMD | 188.1 GB/s |
| csveee DFA | 106.1 GB/s |
| DataFusion | 15.8 GB/s |
| DuckDB | 7.4 GB/s |
| Polars | 4.1 GB/s |
もちろん「普通のSSDからCSVを180GB/sで読める」という意味ではありません。
むしろ多core環境でParser側の処理をmemory bandwidth付近まで押し上げられるほどscalingしたと見る方が分かりやすいと思います。
CSV以外にも使えそう
ここからは論文ではなく、私の考えた話です。
この方法を CSVを高速化する特殊な技術 ではなく、
入力を小さな単位で生成し、その場で処理し、局所stateだけを最後に統合する
という設計として見ると他のformatにも広げられそうです。
例えば、
CSV → row
JSONL → object
Markdown → block
Text → bounded span
PDF → page / text block
DOCX → paragraph / table
Parquet → row group
といった単位で処理できます。
すべてを完全なstreaming parserにする必要はなく、formatごとに適切なbounded unitを決めればよい。
大量ファイルならfile-level parallelism、巨大な単一ファイルならformatが許す範囲でintra-file parallelism、という形にも応用できそうです。
おわりに
今回面白かったのは、 「Parserをもっと速くする」より前に、「Parserがrecordを返すという形で本当にいいのか」を疑ったこと でした。
interfaceを変えることでparse・processing・parallelismをまとめて再設計できる。
こういう設計変更はかなり好きです。
現在はこの方法をCSV以外にも当てはめて試しているので実測できたらまた改めて報告…できたらします。
参考
-
Simon Ellmann, Thomas Neumann, One Pass to Parse Them All: Fused Parallel CSV Processing
https://db.in.tum.de/~ellmann/papers/csveee.pdf -
ackxolotl/csveee
https://github.com/ackxolotl/csveee -
Simon Ellmann氏によるReddit投稿 Fast CSV parsing in Rust
https://www.reddit.com/r/rust/comments/1wp1esi/fast_csv_parsing_in_rust/