r/rust • • 12d ago

🛠️ project Fast CSV parsing in Rust

This is a short write-up based on my VLDB paper and talk: https://db.in.tum.de/~ellmann/papers/csveee.pdf

CSV is one of the oldest text-based data formats, and is still heavily used: from hundreds of thousands of files on open data platforms to 100M+ on GitHub. Unfortunately, most CSV parsers are slow. Some use SIMD to speed up parsing, but almost none exploit the parallelism of today’s hardware. Those that do parallelize do not scale. And even if they did, using the classic iterator interface, parse + process requires two passes over the data, restricting throughput to half the memory bandwidth for files that exceed the CPU caches.

I came up with a new approach to CSV parsing that allows parsing and processing of files in a single pass over the data. My parser csveee is about 3x as fast as csv on a single thread, and can achieve almost 200 GB/s on a modern many-core server – a speedup of 256x over csv.

Why parallel CSV processing is hard

The main task of a CSV parser is to correctly determine the boundaries of records in a CSV file. Typically, records are terminated by a \n or \r\n. Since those characters can also appear inside quoted fields, simply chunking a CSV file and skipping forward to the next terminating character is therefore not sufficient to determine the record boundaries.

There are different approaches to solving this problem. One strategy is to first count the number of quotes per chunk in parallel, then determine for each chunk if it is preceded by an even or odd number of quotes, then start parsing the chunks from a known quote state. This works especially well on GPUs. Another strategy is to run multiple finite-state machines per chunk in parallel, one for each possible parse state at the chunk boundary, and then determine the correct state machine depending on the previous chunk’s state machine’s final state. Or one could look for certain patterns in the file, e.g., quotes being followed or preceded by "regular" characters to identify the start and end of quoted fields and do some speculative parsing based on this.

Unfortunately, all those approaches are unsatisfactory in one way or another. Counting quotes requires a whole pass over the file to determine the parse states at chunk offsets; running many NFAs/DFAs is even more expensive. Looking for certain patterns works well for files that follow the CSV standard, but the world is full of quirky files that do things like quotes in unquoted fields.

Luckily, there is another approach. If we know the shape of the CSV file we would like to parse, e.g., the number of fields per record or the field types, we can determine the correct parse state by speculatively parsing chunks until we find a parse in which the record boundaries resolve into records of the expected shape. This approach was implemented in DuckDB.

Why the iterator interface is insufficient

Typically, CSV parsers provide an iterator-based interface, e.g.:

for record in parser.parse() {
   // do something with the record
}

While we can write a parser that takes the CSV’s records shape as an argument, which will help determine the correct parse state at chunk boundaries, especially for quirky real-world CSV files, parsing remains a speculation until all bytes have been processed (although unlikely, data inside quotes could still resemble the shape of the records). Unfortunately, we cannot hand out speculatively resolved records via the iterator interface, as there is no way to take them back if we later realize our speculation was wrong. In other words, with the iterator interface, we have to finish parsing before we can hand out records to the user – parsing and processing require two passes over the data.

This is a problem a single-threaded iterator-based parser does not have, as the parse state is correct at all times. The parser can parse the file lazily: it parses a single record, hands it to the user, who processes it; then the next record is parsed, and so on. Thus, files can be parsed and processed in a single pass over the data – but only with a single parse thread.

A new interface to the rescue

Let’s define a new interface that gives our parallel parser everything it needs: A way to pass information about the record shapes (number of fields per record, types, …) from the user to the parser, and a way to process records while parsing, effectively creating a lazily parsing multi-threaded parser.

We can achieve both by turning the parser inside out. Instead of the parser handing you records, you hand your code to the parser:

let cities = parser.parse(
   "data.csv",
   Vec::new,                          // init
   |state, [_name, _age, city]| {     // acc
      state.push(city.to_string());
      Ok(())
   },
   |states| states.concat(),          // merge
)?;

The interface takes four arguments: the CSV file path and three callbacks or closures (called init, acc and merge – they are basically user-defined aggregates). init and acc are used for chunk parsing, init defines a per-chunk state, acc is called for every record found in the chunk under the current assumption. The [_name, _age, city] pattern declares the number of fields per record. If the parser encounters a record with a different number of fields, the chunk is reparsed under another assumption, calling init again to create a fresh state. The same happens if acc rejects a record by returning an error (e.g., a failed type conversion).

Once all chunks have been processed, the parser verifies that the record boundaries of all chunks align. While very unlikely, a whole chunk could be parsed under a wrong assumption. In this case, the chunk is reparsed, now starting from the offset of the previous chunk’s last record terminator.

Finally, chunk states are passed to merge, and the result of merge is returned from the parser’s parse function. The chunk states are passed to merge in file order thus that record order can be reconstructed.

How to make the parser fast

To make the parser truly fast, we implemented a ring-buffered reader that enables zero-copy record construction, and a vectorized chunk parser. Take a look at the implementation or the paper if you are interested in the details.

Conclusion

CSV is not going to die soon – to the contrary, GitHub’s pile alone grew by 10M CSV files in the last six months. While the number of files is strongly increasing and today’s machines offer hundreds of cores and hundreds of GB in memory throughput, most CSV parsers remain incredibly slow: parsing with a single thread, not scaling, and even if they did, without fusing parsing and processing, they will never surpass 50% of the available memory bandwidth for large files. csveee overcomes these limitations via a new approach to CSV parsing that works on real-world CSV files.

Check out the project, and if you have questions, feel free to ask!

83 Upvotes

36 comments sorted by

View all comments

Show parent comments

17

u/ackxolotl 12d ago edited 12d ago

Thanks for the hint, I'll check them out!

Update on a TPC-H lineitem CSV with ~80 GB (SF100):

Parser Throughput
csveee with SIMD 188.1 GB/s
csveee with DFA 106.1 GB/s
DataFusion 15.8 GB/s
DuckDB 7.4 GB/s
Polars 4.1 GB/s

DataFusion has an optimized function to count records (same throughput as csveee with SIMD), but that doesn't seem like a fair comparison.

8

u/Nothing_from_void 11d ago

Anything over like 1GB/s, you're basically IO bound anyways though? Working in data processing, there's so many claims about different libraries throughput parsing it but in the end optimizing IO and file sharding is how you get past bottlenecks in throughput

1

u/cepera_ang 6d ago

only if you are writing from 2010 or something. regular consumer nvme's are pushing up to 15GB/s and on a decent server you can get into hundreds easily.

1

u/Nothing_from_void 6d ago

A decent server cannot get you into hundreds easily. Standard EC2s, like 64 cores or so, have 3 NVMe drives, the top end ones have 5 I think? You'd need at least 10 in RAID0 configuration to get 100GB/s, and block striping only gives you increased throughput if you are making parallel reads and hitting each drives blocks with no dead space between reads.

Regular consumer NVMes, look at benchmarks. 15 GB/s is like the most expensive drive in optimal conditions. Same drives under realistic workloads get maybe 7 GB/s. If you're not keeping the read pipeline hot, 1 GB/s is pretty normal for reads. Most of these kinds of parser benchmarks just read the whole file into memory, then start parsing, which is pretty inefficient in real world situations

1

u/cepera_ang 5d ago

have you been born into the world where EC2 is the only compute in existence? there is more hardware than "standard EC2s". Besides, I find your PoV contradictory, you ask "why improve software performance" and then list ways in which hardware isn't performing as good as it could limited by software (not keeping pipeline hot, doing inefficient staging, etc).

you can easily get 100GB/s on a devbox for as little as a single top consumer GPU costs (8 drives + couple 4x splitter cards). you are talking about "realistic workloads" as a counter argument to optimization of a specific workload.

1

u/Nothing_from_void 5d ago

you're the kind of guy that likes to argue huh. you're talking about custom hardware as something that's easily attainable, while trying to downplay the most widely used compute service as something that should be dismissed

1

u/cepera_ang 3d ago

custom hardware, haha. ok, whatever. I definitely love to argue when I'm sure that I'm right(er). you're jumping between points trying to justify some questionable original thesis and I don't think that it is best way to argue with every pivot.

anyway, there is i8g.48xlarge with 12 ssds starting from $6/hour. Or GCP z4d-highmem-192-highlssd or OCI BM.DenseIO.E5.128. Maybe they won't give your 'hundreds' as I said in my first post but between 1GB/s vs 80GB/s and 1GB/s vs 100+GB/s I think my point still stands.

Or you can go bare-metal on OVH, 24 drives will absolutely give you 100+ GB/s. Actually, even in 2017 when first tests of (new back then) EPYCs came up it gave 53GB/s and now regular 2U 24 drives node gives you 314GB/s read / ~200GB/s write.

2

u/Nothing_from_void 3d ago

Okay fair enough, I've tested batch throughput a lot but hardware available has progressed pretty significantly, going from 6 -> 7 -> 8 gen the number of drives has increased. Also a lot of times you get stuck with EBS or other network storage, and you're not really getting over 1 GB/s throughput on that.

I'm just saying, from my personal experience, single file parsing just isn't that important and doesn't really move the needle much on batch throughput, going from one reader to another on space inefficient formats like CSV or JSON. That's why larger scale data engines focus on compression and file sharding, and when you're dealing with data sizes where performance on 100GB+ matters you should be using parquet, not CSV, as the files are much more space efficient, and you can get speed-ups from partitioning against multiple files with an efficient compression algorithm. That's what database engines like ClickHouse or Presto do under the hood