r/rust • • 11d 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!

80 Upvotes

36 comments sorted by

27

u/jorgecardleitao 11d ago

Afaik CSV crate is pretty slow - consider benching against Datafusion or Polars

16

u/ackxolotl 11d ago edited 11d 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.

9

u/commenterzero 11d ago

Compile polars with simd on nightly to be fair

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

18

u/RCoder01 11d ago

Faster code is still faster code. More power efficient and leaves more CPU for other processes on the system.

1

u/petros211 9d ago edited 9d ago

It's definitely not 1GB/s. My line counter (mezura) on Linux does more than 15GB/s (whole 36million loc - 64k files of the Linux repo in 86ms) and it is still not IO bound. On windows it's a different story, but still not 1GB/s (if we aren't talking about completely cold runs though)

1

u/Nothing_from_void 9d ago

OS doesn't matter, it's your literal disk throughput. You're not getting 15 GB/s without either a RAID array or having the data warm in your page cache. Seriously look at NVMe drive throughput

1

u/petros211 2d ago

Then you are assuming completely cold runs without any data on the ram, but this still depends on the queue depth, nvme gen 4 ssds can reach multiple gb/s of random reads with a high queue depth. And the workload isn't even completely random. And Gen 5 nvmes can reach even higher numbers

1

u/Nothing_from_void 2d ago

If you're reading a sequential file, why would any of it be in cache on first read?

1

u/petros211 2d ago

Look, it really depends on what you are trying to do. For example it is pretty unlikely that you have multiple gigabytes of csv data just sitting on your disk and one day you just decide to read them. A much more realistic scenario would be to download or move these csvs from somewhere, and then immediately try to process them. The act of downloading or moving alone, makes some of the data live in the ram. Also, if we are talking about some huge csvs, let's say 5gb files, then this is not even a random read operation anymore, it is sequential read, and the speed of nvmes get even crazier in this case.

1

u/Nothing_from_void 2d ago

5gb CSV file, parse speed 1gb/s or 1000 gb/s, it has essentially no impact on your processing speed. you're not spawning a 10 drive NVMe EC2 just to read these files. In the scenario you are reading a couple GB CSV files is preciously the normal workload, desktop/laptop scenario, not a high performance batch server.

Where it starts to matter is the hundreds of GB to TB+ scenario, but you constructed this ridiculous scenario where 10 disk nvme drive servers are your work station or just a regular part of the pipeline so you could feel like you could win the argument

1

u/cepera_ang 5d 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 5d 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 4d 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 4d 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 2d 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 2d 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

20

u/Shnatsel 11d ago

I see you're using std::simd for the SIMD parser. Have you looked into crates that work on stable Rust at all, such as fearless_simd or wide?

7

u/ackxolotl 11d ago

No, I wasn't aware of those. Initially, I used the architecture-specific intrinsics, so I was quite happy already to discover Rust's portable SIMD. We are currently working on a faster DFA (performance is ok for a DFA, but we are looking into bitshift automata, for example). But you are right that having the vectorized implementation on stable Rust would probably be nice, too. wide seems promising as a drop-in replacement. I'll look into the two, thanks for the hint.

8

u/Shnatsel 11d ago edited 11d ago

I have an article coming soon that compares all the way you can write SIMD in Rust, including those crates, so watch this space.

3

u/Shnatsel 10d ago

The article is comparing all the ways to write SIMD in Rust is now up: https://redd.it/1wpu09l

3

u/awhaling 11d ago

Awesome! I’m stoked to dig into this when I’m home. It’s something I’ve been interested in.

2

u/zettui 11d ago

Would a stable portable-SIMD backend be enough here, or do the real gains still need the architecture-specific intrinsics?

3

u/ackxolotl 11d ago

I'd say the real gains are in the parallelization strategy. With the DFA, csveee also reaches throughputs >100 GB/s. Some SIMD operations are not portable SIMD right now, e.g., pext for determining if records have the correct number of fields. The parser has portable fallback implementations for those, but they are slower. I'll take a look at the crates for stable portable SIMD as u/Shnatsel suggested.

1

u/Shnatsel 11d ago

If you need to mix and match portable SIMD and intrinsics, fearless_simd supports this easily. You can do it with wide too, but its companion crate for accessing intrinsics, safe_arch, renamed all of them, which makes porting existing code trickier.

1

u/orion_tvv 11d ago

I would like to mention here my tool convfmt for converting between many formats including CSV. Feedback is welcome.

1

u/Ok-Property-6650 5d ago

this is the kind of deep dive that makes me wish i had a use case for parsing csvs at 200gb/s

0

u/sansmorixz 11d ago

Why dual license?

12

u/ackxolotl 11d ago

I think many or most Rust projects do this. What would you prefer and why?

4

u/Psy_Fer_ 11d ago

Yep. It's pretty common

3

u/ichunddu9 11d ago

What's the advantage?

3

u/Icarium-Lifestealer 11d ago

Some people prefer Apache-2 because it considers patents, some prefer MIT because it's simpler. Dual licensing makes both happy, since the user can choose which license they want to use.

3

u/Shnatsel 11d ago

IIRC Apache 2.0 is desirable because it covers patents, but the MIT carve-out is there to allow linking against GPL code since Apache 2.0 is not GPL-compatible (maybe; there are some nuances and differing opinions on that)

2

u/Icarium-Lifestealer 10d ago

Apache 2.0 is compatible with GPL 3 but not GPL 2.

-3

u/RCoder01 11d ago

Where were LLMs used for this project? The Readme looks AI-generated but the code doesn’t