The One Billion Row Challenge in CUDA
tspeterkim.github.io
tspeterkim.github.io
After you deal with parsing and hashes, basically you are IO limited, so mmap helps. The C code takes less than 1.4s without any CUDA access. Because there is no compute to speak of, other than parsing and a hashmap, a reasonable guess is that even for the optimal CUDA implementation, the starting of kernels and transfer of data to the GPU would likely add a noticeable bottleneck and make the optimal CUDA code slower than this pure C code.
If you account for time loading from disk, the C implementation would be more like ~5s as reported in the blog post [1]. Speculating that their laptop's SSD may be in the 3GB/s range, perhaps there is another second or so of optimization left there (which would roughly work out to the 1.4s in-memory time).
Because you have a lot of variable width row reads this will be more difficult on a GPU than CPU.
Let's say I create a "cache" where I store the min/mean/max output for each city, mmap it, and read it at least once to make sure it is in RAM. If the cache is available I simply write it to standard out. I use whatever method to compute the first run, and I persist it to disk and then mmap it. The first run could take 20 hours and gets discarded.
By technicality it might fit the rules of the original request but it isn't an interesting solution. Feel free to submit it :)
This is a fun project for learning CUDA and I enjoyed reading about it——I just wanted to point out that the performance tuning in this case is really on the parsing, hashing, memory transfers, and IO. Taking IO out of the picture using specialized hardware or Linux kernel caching still leaves an interesting problem to solve and the focus should be on minimizing the memory transfers, parsing, and hashing.
[0] https://github.com/gunnarmorling/1brc/issues/189#issuecommen...
on large datasets, once loaded into GPU memory, cross GPU shuffling with NVLink is going to be much faster than CPU to RAM.
on the H100 boxes with 8x400Gbps, IO with GDS is also pretty fast.
for truly IObound tasks I think a lot of GPUs beats almost anything :-)
https://docs.nvidia.com/gpudirect-storage/overview-guide/ind...
You're yada-yada-yadaing the best part.
If the disk can process the data at 1GB/s, the CPU can process the data at 2GB/s, the GPU can process the data at 32GB/s, then the CPU can process the data at 1GB/s and the GPU can process the data at 1GB/s.
(also, personally, "large dataset" is a short way to say "doesn't fit in memory". if it fits in memory it's small, if it doesn't fit in memory it's large. but that's just my opinion. I generally avoid calling something a "large" or "small" dataset because it's an overloaded term that means different things to different people.)
This problem is quite a Nerd Snipe so it a good thing that I dont have a computer with more than 16 GB or I might end up trying it myself.
This is similar to how GPU acceleration is bolted on in a ton of commercial software, you look at hotspots and then move just those over to the accelerator. For example outsourcing some linear algebra in a FEM solver. That's a sensible approach for trying to accelerate legacy software, but it's not very effective and leads to poor utilization ("copy mat, copy mat, copy mat copy copy copy mat" - Dave Cutler about GPU acceleration). I think a) you can do all of this on the GPU b) this should be limited by PCIe bandwidth for getting the file contents into VRAM. The 1BRC data set is 14 GB. So this is processing that file at about 1 GB/s using both CPU and GPU.
The catch is that this map must fit in shared memory, which is pretty limited on all current hardware: ~100KB.
I originally thought that my map (stats array) was too big to fit into this shared memory. Now, however, I realize it can. It'll interesting to see how much speedup (or not!) this optimization can bring.
In many cases it is better to just pass the raw byte offsets to the workers, then they can skip to the first newline and progress until the next newline after their "end offset". This way the launcher doesn't need to access the input data at all.
This can be ineffective when boundaries are not self-describing or when the minimum processable chunk is near the expected batch size (as you will have significant swings in batch size) but I think would work quite well for this case.
I also already pass these offsets to the threads as `Part* parts`.
I also probably didn't understand your suggestion and am drawing a blank here. So pls feel free to elaborate and correct me.
for (int i = 0; i < parts[bx].length; i++) { // bx is the global thread index
char c = buffer[parts[bx].offset-buffer_offset + i];
by something like long long split_size = size / num_parts;
long long offset = bx * split_size;
while (buffer[offset++] != '\n')
;
for (int i = 0; buffer[offset+i] != '\n'; i++) {
char c = buffer[offset + i];
(ignoring how to deal with buffer_offset for now). for (int i = 0; buffer[offset+i] != '\n'; i++) {
This would only process the current line, though. Here, each thread processes ~split_size bytes (multiple lines).Even if were to read multiple lines, how would a thread know when to stop? (at which offset?)
And when it does, it should communicate where it stopped with the other threads to prevent re-reading the buffer. My brain's hurting now.
let chunk_size = buf.length / threads;
let start_byte = chunk_size * bx;
let end_byte = start_byte + chunk_size;
let i = start_byte;
while (buffer[i++] != '\n') {} // Skip until the first new row.
while (i < end_byte) {
let row_start = i;
let row_end = i;
while (buffer[row_end++] != '\n') {}
process_row(buffer, row_start, row_end);
i = row_end + 1;
}
Basically you first roughly split the data with simple byte division. Then each thread aligns itself to the underling data chunks. This alignment can be done in parallel across all threads rather than being part of a serial step that examines every byte before the parallel work starts. You need to take care that the alignment each thread does doesn't skip or duplicate any rows, but for simple data formats like this I don't think that should be a major difficulty. while (i < end_byte) {
Comparing it to my original solution, 50X divergent branches are introduced! (ncu profiling)The only difference between the two is that the for loop could deterministically iterate, yet this while loop iterates for an unknown amount (at kernel launch time).
I admit, I don't perfectly understand the reason. But this is the most likely culprit.
The same level of performance can be obtained using an Ada like an L40s or and RTX 4090.
The transfer across the nvlink connecting the Cpu and GPU on a gh-200 after the parse and encoding of the source CSV takes a negligible amount of time given the 500 gb/sec of system memory bandwidth and the 900 GB /sec interconnection between Cpu and Gpu.
So the problem is disk bandwidth that's going to limit the performance of the Gpu kernel. The faster solution should be parse the Csv with a gpu kernel using namp and Managed memory (?) encode the station into a interger or a small integer. The min and max value can be used to create keyless perfect hash table for each SM to limit the concurrency on global memory using 32bit atomic operations for min, max, count and sum and then do a final reduction on the Gpu.
I don't think that is needed more then 1 modern gpu for this, especially if you are on a modern hardware like the gh-200.
I'm running this kind of aggregates on a gh-200 using 10 Billion of records having the data in Gpu or Cpu memory using our software (heavydb) for testing purposes in the last two weeks
[0] https://github.com/gunnarmorling/1brc#evaluating-results
Do you think the L4 should have similar perf? Considering it has similar features and specs to L40, and I doubt the lower 300 GB/s bandwidth woudl be an issue.
Something else to consider is cost of cloud hardware, though it could be extrapolated from the author's results.
For this problem I think coursening would also help (so everything isn’t queued up on the copy to global memory)
By coarsening, do you mean making the threads handle more file parts, and reducing the number of private copies (of histogram or stats here) to globally commit at the end?
The same level of performance can be obtained using an Ada like an L40s or and RTX 4090.
The transfer across the nvlink connecting the Cpu and GPU on a gh-200 after the parse and encoding of the source CSV takes a negligible amount of time given the 500 gb/sec of system memory bandwidth and the 900 GB /sec interconnection between Cpu and Gpu.
So the problem is disk bandwidth that's going to limit the performance of the Gpu kernel. The faster solution should be parse the Csv with a gpu kernel using namp and Managed memory (?) encode the station into a interger or a small integer. The min and max value can be used to create keyless perfect hash table for each SM to limit the concurrency on global memory using 32bit atomic operations for min, max, count and sum and then do a final reduction on the Gpu.
I don't think that is needed more then 1 modern gpu for this, especially if you are on a modern hardware like the gh-200.
I'm running this kind of aggregates on a gh-200 using 10 Billion of records having the data in Gpu or Cpu memory using our software (heavydb) for testing purposes in the last two weeks
How bad would performance have suffered if you sha256'd the lines to build the map? I'm going to guess "badly"?
Maybe something like this in CUDA: https://github.com/Cyan4973/xxHash ?
However, the problem of collisions across threads and dealing with concurrent map key insertions still remains. e.g. when two different cities produce the same hash (one at each thread), how can we atomically compare 100 byte city strings and correctly do collision-solving (using linear probe, for example - https://nosferalatu.com/SimpleGPUHashTable.html)
Atomic operations are limited to 32-bits.
I'm using 64bit atomics at work, are you on an old version of cuda or are some operations only supported on 32bit?
Atomics work up to 128-bits (https://docs.nvidia.com/cuda/cuda-c-programming-guide/#atomi...).
Regardless, it's still less than 100 bytes, which is the max length of city strings.
My hope with this blog was to rile up other CUDA enthusiasts. Making them wanna bring in bigger and better hardware that I don't have access to.
+ Someone in the CUDA MODE community got 6 seconds on a 4090.
I wonder if one could use a reduce operation across the cities or similar rather than atomic operations? GPUs are amazing at reduce operations.
The problem was that the work of gathering all the temperatures for each city (before I could launch the reduction CUDA kernels) required a full parsing through the input data.
My final solution would be slower than the C++ baseline since the baseline already does the full parsing anyways.
can you not stop early if extracted value becomes smaller than val?
so just using cuda's AtomicMin for ints would be enough