The One Billion Row Challenge
morling.dev
morling.dev
[0] https://github.com/gunnarmorling/1brc/blob/main/src/main/jav...
[0] https://twitter.com/mtopolnik/status/1742652716919251052
That might help against algorithms that are (accidentally, or on purpose) tuned to the specific dataset.
It's very cheap to detect station name collisions - just sample ~10,000 points in the file, and hope to find at least one of each station. If you find less stations than the full run, you haven't yet seen them all, keep hunting. If you find two that collide, fall back to the slow method. Checking 0.00001% of the dataset is cheap.
(Describing only the 'happy path' here - other paths can be made fast too, but will require different implementations)
* Since temperatures are only to 0.1 decimal points, we have a finite number of temperatures. ~400 temperatures will cover all common cases.
* We also have a finite number of place names. (~400)
* Just make a lookup table of all temps and all place names. (~160,000)
* autogenerate a state machine that will map each of the above 160,000 things, at any rotation within a 4 byte register, to a unique bin in a hash table. The state machine will have one 32 bit state register (16 bits to output, 16 to carry over to the next cycle) where every cycle a lookup in the state transition table is done, and the next 4 bytes of data XOR'ed on top.
* run through all the data, at RAM speed, incrementing counters for each state the machine ends up in. (there will only be 65k). These counters fit fully in cache.
* With just 32 bits of state, with AVX512 we can be running 512 copies of this in parallel if we like, per core! AKA, compute will not be the bottleneck.
* from the values of the counters, you can calculate the answers. 65k is a far smaller number than 1 billion, so you don't need to do this bit fast.
* for anything that doesn't map to a valid bin (higher/lower temperatures, unknown place names), just fallback to slow code (one state of the state machine can be reserved for 'escape to slow code'. Use these escapes for min/max handling too, since it will only happen a few thousand times).
* I think this approach can operate at RAM speed with just one core with AVX512, so no benefit in splitting across cores.
IO will dominate the running time for this, and JSON parsing will be second.
I always worried it was too easy, but I'm heartened by how many comments miss the insight I always looked for and you named: you don't need to store a single thing.
I do wonder if it'll work here, at least as simply as the tracking vars you mention, with so many rows, and the implicit spirit of "you should be able to handle _all_ data", overflow might become a legitimate concern - ex. we might not be able to track mean as simply as maintaining sum + count vars.
* "oh sorry times up, but the right answer is a red black self balancing terenary tree recursive breadth-first search", ironically, was one of my Apple interviews.
(Here there is a restricted range for the temperature with only 199 possible values (-99.9 to 99.9 with 0.1 increment) so you could do it constant memory, need something like 4*199 bytes per unique place name))
For the sum overflow is not an issue if you use 64-bit integers. Parse everything to integers in tenths of degree and even if all 1 billion rows are 99.9 temperature for same place name (worst possible casE), you are very far from overflowing.
I am silly and wrote 'mode' and shouldn't have :P (wetware error: saw list of 3 items corresponding to leetcode and temperature dataset, my 3 were min/max/average, their 3 are mean/median/mode)
All far more expensive than the state machine approach which is pretty much 'alignment doesn't matter, ascii doesn't matter, we're gonna just handle all possible states that could end up in a 32 bit register').
If you came up with them on the fly, then... well, sign here for your bonus. Can you start Monday?
It's not "non-trivial" it's impossible. I'm not sure why people think median can be approximated at all. You need to look at every data point and store a counter for the lesser of: (a) all possible values, or (b) all elements. Just consider a data set with 1 million ones and 999,999 zeroes. Throwing away (or failing to look at) one single number can give you an error of 50%, or 100% for two numbers. If you want to make it really nasty, throw in another million random 64-bit floats between [-0.1, 0) and a million between (1, 1.1]. Four million elements, two million of them unique, in the range [-0.1, 1.1], and failing to account for two elements can get your answer wrong by +/- 1.0.
Unless you start by sorting the list, but I wouldn't call that a "streaming" calculation.
If someone on HN knows how to do it, they will jump in and tell me exactly why it's "trivial," and I'll get some nice R&D work for free. Of course, it will probably involve mounting an FTP account, spinning up a CA and running rsync...
n.b. to you and your parent comment's commenter: you're the only two who used the word trivial in the entire comments section modulo another comment not in this thread that contains "SQL in theory makes this trivial...". The HNer claiming to build a streaming median function in a weekend fantasy might not come to fruition :(
A quick search shows that Boost.Accumulators have several median estimation implementations available to choose from: https://www.boost.org/doc/libs/1_84_0/doc/html/accumulators/...
> Median estimation based on the P^2 quantile estimator, the density estimator, or the P^2 cumulative distribution estimator.
But it brings up an important point: real data that you actually encounter is not a random sampling from "the general case". It is often possible to do a pretty good job approximating the median from real-world data, if you have some understanding of the distribution of that data. You just have to accept the fact that you might be totally wrong, if the data behaves in ways you don't expect. But of course, the more you have to know about your data before running the algorithm, the less useful the algorithm itself is.
The difference is whether you can make guarantees or not. Most algorithm design is concerned about such guarantees: on all possible inputs, this algorithm has at worst such-and-such performance. Hence the reliance on Big-O. You cannot ever make such guarantees on a "median approximator" without specifying some sort of precondition on the inputs.
What guarantees do we want to make? With "an estimator", it'd be nice to say something like: we'd like our approximation to get better given more values. That is: the more values we see, the less the next single value should be able to change our approximation. If you've looked at a billion values, all between [min, max], it'd be nice if you knew that looking at the next one or two values could only have an effect of at most 1 / f(1 billion) for some monotonically increasing f. Median does not have that property: looking at just two more values (even still within the range of [min, max]) could move your final answer all the way from max to min. If you stopped 2 data points earlier, your answer would be as wrong as possible. This remains true, for some inputs, even if you've looked at 10^10^10^10^300 values. The next 2 might change your answer by 100%.
Your state machine would need at least 2*160,000 states (you need an extra bit to flag whether you have reached a newline in the last word and need to increment a counter or not), correct? And you are assuming the input is 4 bytes, so won't your transition table need (2^32)*2*160,000 ~ 10^15 entries (at least 3 bytes each)?
You could shrink your input size to 2 bytes but then you can't work on a word at a time, and for a realistic number of relevant states your transition table is still way bigger than you can fit in even L3 cache.
Unless I am missing something very basic, this doesn't seem like a viable approach.
Memory bandwidth might dominate, but probably not I/O. The input file is ~12GB, the machine being used for the test has 32GB, and the fastest and slowest of five runs are discarded. The slowest run will usually be the first run (if the file is not already cached in memory), after which there should be little or no file I/O.
What if the code was run against constraints such as Max memory limit in a docker container.
You can do this whole thing in about 24 seconds, as long as you're smart about how you chunk up the text and leverage threading. https://github.com/coriolinus/1brc/blob/main/src/main.rs
Having a quick look at your code, couple of thoughts:
- You shouldn't bother with parsing and validating UTF-8. Just pretend it's ASCII. Non ASCII characters are only going to show up in the station name anyway, and all you are doing with it is hashing it and copying it.
- You are first chopping the file into line chunks and then parsing the line. You can do it it one go, just look at each character byte by byte until you hit a semicolon, and compute a running hash byte by byte. You can also parse the number into an int (ignoring decimal point) using custom code and be faster than the generic float parser.
- If instead of reading the file using standard library, you mmap it, should also speed things up a bit.A single core can't saturate RAM bandwidth. Cores are limited by memory parallelism and latency. Most modern x86 server chips can retire 2 SIMD loads per clock cycle (i.e. 32GB/s @ 1GHz using AVX2), so AVX-512 is not necessary to max out per-core bandwidth - but if you're loading from DRAM you'll likely cap out much, much earlier (closer to 10-16GB/s on commodity servers).
As long as the majority of your data spills into RAM, you'll see a massive throughput cliff on a single core. With large streaming use-cases, you'll almost always benefit from multi-core parallelism.
It's very easy to test: allocate a large block of memory (ideally orders of magnitude larger than your L3 cache), pre-fault it, so you're not stalling on page faults, and do an unrolled vector load (AVX2/AVX-512) in a tight loop.
I was also intrigued by the same statement and was wondering about it as well. For AMD I don't know but on Intel I believe that each L1 cache-miss is going to be routed through the line-fill-buffers (LFB) first before going to L2/L3. Since this LFB is only 10-12 entries large, this number practically becomes the theoretical limit on the number of (load) memory operations each core can issue at a time.
> Most modern x86 server chips can retire 2 SIMD loads per clock cycle (i.e. 32GB/s @ 1GHz using AVX2),
I think in general yes but it appears that ALU on Alderlake and Saphire Rapids has 3 SIMD load ports so they can execute up to 3x SIMD AVX2 mm256_load_si256 per clock cycle.
But for that, you're right, the workload must be split across as many cores as necessary to keep DRAM maxed out or you're throwing away performance. Luckily this problem is easily split.
How can the state machine be run in parallel, when the next state always has a dependency on the previous state?
Also, how exactly would the state register be decoded? After you XOR it with 4 bytes of input, it could be practically any of the 4.7 billion possible values, in the case of an unexpected place name.
And even for expected place names longer than 4 bytes, wouldn't they need several states each, to be properly distinguished from other names with a common prefix?
> Q: Can I make assumptions on the names of the weather stations showing up in the data set?
> A: No, while only a fixed set of station names is used by the data set generator, any solution should work with arbitrary UTF-8 station names (for the sake of simplicity, names are guaranteed to contain no ; character).
You'd have an upper limit on how many place names you can have. But no hard coding of what they are.
You are basically describing a perfect hash, right?
Remember that in the extreme, a state could mean "city A, temp zero, city b" - "0\rB" so it isn't a direct mapping from state machine state to city.
The state machine generation would need to know all 400 - but it's easy enough to scan the first few hundred thousand rows to find them all, and have a fallback incase a name is seen that has never been seen before. (the fallback would be done by having the state machine jump to an 'invalid' state, and then at the end you check that that invalid state's counter is zero, and if it isn't, you redo everything with slow code).
The README says the maximum number of unique place names is 10 000, so you should probably design for the fast case to handle that many instead of just going by what is in the database that it happens to be measuring with at the moment.
But more seriously, the JVM's support for vector intrinsics is very basic right now, and I think I'd spend far more time battling the JVM to output the code I want it to output than is fun. Java just isn't the right language if you need to superoptimize stuff.
Theoretically all of the above is super simple SIMD stuff, but I have a suspicion that SIMD scatter/gather (needed for the state lookup table) isn't implemented.
I still appreciate the intention of seeing how far we can squeeze performance of modern Java+JVM though. Too bad very few folks have the freedom to apply that at their day jobs though.
For more such tasks you could go to highload.fun.
Is the data given in RAM or do you have to wait for spinning rust I/0 for some GB amount of data?
I think it's better to discard the two slowest, or simply accept the fastest as the correct. There's (in my opinion) no good reason to discard the best runs.
If you are looking for a real-world, whole-system benchmark (like a database or app server), then taking the average makes sense.
If you are benchmarking an individual algorithm or program and its optimisations, then taking the fastest run makes sense - that was the run with least external interference. The only exception might be if you want to benchmark with cold caches, but then you need to reset these carefully between runs as well.
If the language is garbage collected, or if the test is randomized you obviously don't want to look at the minimum.
Depends what you’re estimating. The minimum is usually not representative of “real world” performance, which is why we use measures of central tendency over many runs for performance benchmarks.
https://tratt.net/laurie/blog/2019/minimum_times_tend_to_mis... is a good piece on the subject.
There can't be outliers in terms of fastness; the CPU doesn't accidentally run the program much faster.
But then again, what the hell do I know...
> write a Java program for retrieving temperature measurement values from a text file and calculating the min, mean, and max temperature per weather station
Depending on how far you want to stretch it, doesn’t precomputing the result on 1st run count too? Or pre-parsing the numbers into a compact format you can just slurp directly into rolling sums for subsequent runs.
Not the slightest in the spirit of the competition, but not against the rules as far as I can tell.
Edit: if we don't like pre-computation, we can still play with fancy out-of-the-box tricks like pre-sorting the input, pre-parsing, compacting, aligning everything, etc.
Edit 2: while we're here, why not just patch the calculate_time script to return 0 seconds :) And return 9999 for competitors, for good measure
This might become a contest of judging what is fair pre-computing and what is not.
That's why machine learning contests don't let participants see the final data.
A run consists of:
* I generate random measurements with the provided script, which will be stored in a readonly file.
* I run your program to calculate the results (and time it).
* I run my baseline program and your entry is valid if my baseline program calculates the same results.> The computation must happen at application runtime, i.e. you cannot process the measurements file at build time (for instance, when using GraalVM) and just bake the result into the binary
The first run is indeed runtime, not build time, so technically I think I still count with pre-computation, sending that to the daemon (or just stashing it somewhere in /tmp, or even crazier patching the jarfile), and just printing it back out on the second run onwards.
Why not preprocess it into a file that contains the answer? (-: My read of the description was just: you're meant to read the file as part of the run.
It would also depend of how many different stations there are and what they are for the hash lookup, but seriously I doubt this will be anything measurable compared to I/O.
Disk access can certainly be parallelized, and NVMe is blazing fast, so the bottleneck is more the CPU than disk. There are systems that are built around around modern hardware that realize this (redpanda.com, where I work is one such example)
Parsing is a lot of the compute time and some SIMD tricks like SWAR for finding delimiters can be helpful. Stringzilla is a cool library if you want to see a clean implementation of these algorithms: https://github.com/ashvardanian/StringZilla
See my reply here about the file being cached completely in memory after the first run: https://news.ycombinator.com/item?id=38864034
The AMD EPYC-Milan in the test server supports memory reads at 150 gigabytes/sec, but thats a 32 core machine, and our test only gets 8 of those cores, so we probably can't expect more than 37 gigabytes per second of read bandwidth.
The total file is ~12 gigabytes, so we should expect to be able to get the task done in 0.3 seconds or so.
There are only ~400 weather stations, so all the results fit in cache. It's unclear if it might be faster to avoid converting ascii to a float - it might be faster to add up all the ascii numbers, and figure out how to convert to a float only when you have the total.
Saturating NVMe bandwidth is wicked hard if you don't know what you're doing. Most people think they're I/O bound because profiling hinted at some read() function somewhere, when in fact they're CPU bound in a runtime layer they don't understand.
How? EPYC-Milan features PCIe 4.0 and theoretical maximum throughput of x4 sequential read is ~8GB/s.
In bare-metal machines, this figure is rather closer to 5-6GB/s, and in cloud environments such as CCX33, I imagine it could be even less.
So, unless I am missing something you'd need ~2 seconds only to read the contents of the ~12GB file.
To be fair, with only additions, and only a billion of them, you're probably fine working in floating point and rounding to one decimal in the end.
Anyway the OP's file only has 1B rows so it'll be ~1GB compressed which fits in memory after the first discarded run. IO bandwidth shouldn't matter in this scenario:
> All submissions will be evaluated by running the program on a Hetzner Cloud CCX33 instance (8 dedicated vCPU, 32 GB RAM). The time program is used for measuring execution times, i.e. end-to-end times are measured. Each contender will be run five times in a row. The slowest and the fastest runs are discarded
Edit: The github repo says it's 12 GB uncompressed. That confirms that IO bandwidth doesn't matter
Depends on the OS and filesystem being used. The input file is ~12GB, and it's being run 5 times on a machine with 32GB, so after the first run the file could be cached entirely in memory. For example, Linux using ext2 would probably cache the entire file after the first run, but using ZFS it probably would not.
<< Just clarified this in the README: * Input value ranges are as follows:
- Station name: non null UTF-8 string of min length 1 character and max length 100 characters
- Temperature value: non null double between -99.9 (inclusive) and 99.9 (inclusive), always with one fractional digit >>
If we assume worst case entries: -99.9;<100 char station name>
That is 106 chars plus newline (1 or 2 chars). That could be 108 GB of data! And, I want to be so picky here: The README says "UTF-8" and "max length 100 characters". Eh, we could start a holy war here arguing about what does a character mean in UTF-8. Let's be generous and assume it means one (8-bit) byte (like C). I have good faith in the original author!
Some Googling tells me that the current fastest PCIe Gen5 NVMe drive is Crucial T700 w/ Sequential read (max., MB/s) @ 12,400. That is 8.7 seconds to read 108 GB of data. A few other comments mentioned that the sample file is only 12 GB, so that is 1 second, and 100% can be cached by Linux kernel for 2nd and later reads. However, if the test file is greater than 32 GB, you are definitely screwed for file system caching. You will have to pay a multi-second "I/O tax" to read the file each time. This seems to far outweigh any computation time.
Thoughts?
Last, I still think this is a great challenge. I think I just found my next interview question. This question is easy to understand and program a basic implementation. Further discussion can be fractally complex about the myriad of optimizations that could be applied.
No trolling: Can you imagine if someone tries this exercise in Python? There are a goofy amount of optimizations to deal with huge data in that language thanks to SciPy, Pandas and other friends. I am sure someone can write an insanely performant solution. (I know, I know: The original challenge says no external dependencies.)
I have tried 2 naive awk implementations and both took 10x the basic java implementation btw (I'm sure it too can be optimized).
> A: No, while only a fixed set of station names is used by the data set generator, any solution should work with arbitrary UTF-8 station names (for the sake of simplicity, names are guaranteed to contain no `;` character).
I'm unsure if it's intentional or not, but this essentially means that the submission should be correct for all inputs, but can and probably should be tuned for the particular input regenerated by `create_measurements.sh`. I can imagine submissions with a perfect hash function tuned for given set of stations, for example.
This sort of thing shows up all the time in the real world - where, for example, 90% of traffic will hit one endpoint. Or 90% of a database is one specific table. Discovering and microoptimizing for the common case is an important skill.
I think having the names for performance tests public is fine. But there should be correctness tests on sets with random names.
In this case, we can assume place names that can be arbitrary strings and temperatures that occur naturally on earth. Optimizing around those constraints should be fine. We could probably limit the string length to something like 100 characters or so.
Bad assumptions would for example be to assume that all possible place names are found within the example data.
The easiest way to handle Unicode is to not handle it at all, and just shove it down the line. This is often even correct, as long as you don't need to do any string operations on it.
If the author wanted to play Unicode games, requiring normalization would be the way to go, but that would turn this into a very different challenge.
The same logic applies to newline - therefore, you can jump into the middle of the file anywhere and guarantee to be able to synchronize.
I'm no UTF-8 guru, but I think you might be possible to do this sort of a springboard for skipping over multi-byte codepoints, since as far as I understand the upper bits of the first byte encodes the length:
byte utfByte1 = (byte) (val & 0xF0);
if (utfByte1 == (byte) 0xF0) { // 4 byte codepoint
// ignore 3
}
else if (utfByte1 == (byte) 0xE0) { // 3 byte codepoint
// ignore 2
}
else if (utfByte1 == (byte) 0xC0) { // 2 byte codepoint
// ignore 1
} awk -F';' '{
station = $1
temperature = $2
sum[station] += temperature
count[station]++
if (temperature < min[station] || count[station] == 1) {
min[station] = temperature
}
if (temperature > max[station] || count[station] == 1) {
max[station] = temperature
}
}
END {
for (s in sum) {
mean = sum[s] / count[s]
printf "{%s=%.1f/%.1f/%.1f", s, min[s], mean, max[s]
printf (s == PROCINFO["sorted_in"][length(PROCINFO["sorted_in"])] ? "}\n" : ", ")
}
}' measurement.txt CREATE EXTENSION file_fdw;
CREATE SERVER stations FOREIGN DATA WRAPPER file_fdw;
CREATE FOREIGN TABLE records (
station_name text,
temperature float
) SERVER stations OPTIONS (filename 'path/to/file.csv', format 'csv', delimiter ';');
SELECT station_name, MIN(temperature) AS temp_min, AVG(temperature) AS temp_mean, MAX(temperature) AS temp_max
FROM records
GROUP BY station_name
ORDER BY station_name;Very nice! Is this time for the first run or second? Is there a big difference between the first and second run, please?
\copy TEST(CITY, TEMPERATURE) FROM 'measurements.txt' DELIMITER ';' CSV;In the first one, the table is dropped, recreated, populated and queries. In the second example, the table is created from a file FDW to the CSV file. In both examples the loading time is included in the total time
Why not have the machine consider the data it is currently seeing (type, even actual values), think about what end-to-end operation is required, how often it needs to be repeated, make a time estimate (then verify the estimate, change it for the next run if needed, keep a history for future needs), choose one of methods it has at its disposal (index autogeneration, conversion of raw data, denormalization, efficient allocation of memory hierarchy, ...). Yeah, I'm not focusing on this specific one billion rows challenge but rather what computers today should be able to do for us.
time clickhouse local -q "SELECT concat('{', arrayStringConcat(groupArray(v), ', '), '}')
FROM
(
SELECT concat(station, '=', min(t), '/', max(t), '/', avg(t)) AS v
FROM file('measurements.txt', 'CSV', 'station String, t Float32')
GROUP BY station
ORDER BY station ASC
)
SETTINGS format_csv_delimiter = ';', max_threads = 8" >/dev/null
real 0m15.201s
user 2m16.124s
sys 0m2.351
Most of the time is spent parsing the filehttps://stats.stackexchange.com/a/235151/1036
I am not sure if `n*old_mean` is a good idea. Wellford's is typically something like inside the loop
count += 1; delta = current - mean; mean += delta/count;
Even more interesting would be small ram, little local storage and a large file only available via network, I would like to see something other than http but realistically it would be http.
(Of course, the five times in a row might mess with that.)
Parallelism with edge effects is pretty common. Weather simulation, finite element analysis, and big-world games all have that issue. The middle of each cell is local, but you have to talk to the neighbor cells a little.
Code: https://github.com/rockwotj/1brc
Discussion on Filesystem cache: https://x.com/rockwotj/status/1742168024776430041?s=20
In case you haven't noticed yet, the input format guarantees exactly one fractional digit, so you can read a single signed integer followed by `.` and one digit instead.
(The easiest I can think of is submitting reads for "overlapped" chunks, I'm not sure there is an easier way and I'm not sure of how much performance overhead there is to it.)
I thought memory mapping solved a different problem.
I guess you could from Java itself write a new binary and then run that binary, but it would be against the spirit of the challenge.
A fair comparison between languages should include the make and build times. I haven't used Java / Maven for years, and I'm reminded why, heading into minute 2 of downloads for './mvnw clean verify'.
(Also, gradle is faster as a build tool for incremental compilation)
I have no idea what gradle and its daemon do except spending orders of magnitude more time starting up than running the actual build.
You're much better off running javac directly. Then it's fast.
For big projects, compiling only what’s necessary is a huge time win.
Then someone in over 15 years and 8 versions gradle failed to make a system that doesn't require "learning to write sane gradle files" to make sure that it's only two orders of magnitude and not five orders of magnitude slower than it can be.
> and are putting imperative stuff into the configuration phase
Has nothing to do with the slowness that is gradle.
> For big projects, compiling only what’s necessary is a huge time win.
It is, and no one is arguing with that
- in 15 years is still slow by default
- in 15 years could not acquire any sensible defaults that shouldn't need extra work to configure
- is written in an esoteric language with next to zero tools to debug and introspect
- randomly changes and moves thing around every couple of years for no discernible reason and making zero impact on reducing its slowness
- has no configuration specification to speak of because the config is written in that same esoteric programming language
And its zealots blame its failures on its users
Gradle with a daemon is also pretty fast, you are just probably used to some complex project with hundreds of dependencies and compare it to a cargo file with a single dependency.
Not really when it has to run on a cold JVM. And before it warms up, it’s already done.
Then you are in territory of keeping the compiler process between the runs, but that is a memory hog and also not always realiable (gradle often decides to run a fresh daemon for whatever reason).
So theoretically, in lab conditions yes, in practice no.
> while rust is significantly slower
Everybody repeats that but there is surprisingly little evidence in form of benchmarks. I can see rustic compiles over 50k Loc per second on average on my M2, which is roughly the same order of magnitude as Java. And checking without building is faster.
Wiki citation for good measures https://en.wikipedia.org/wiki/Javac.
Maybe you're thinking of Hotspot, which is written in C++.
I've yet to see a project where Gradle daemon a) does anything useful and b) is acutally used by gradle itself (instead of seemingly doing everything from scratch, no idea what it does in the seconds it takes for it to start up).
Also my experience.
Only using Gradle deamon does not help making a multi-project setup much faster. You also have to enable build cache, configuration cache and configure on demand with `--configuration-cache --configure-on-demand` and hope nobody in the project breaks the ability for Gradle to use these caches. But then it still took at least 10 seconds to build and start my services (and that's with incremental builds, like you changed one line of code after the first slow build). I spend two days and more after release to speed this stuff up, before it was 30 seconds sometimes 60 seconds.
And the protobuf Gradle plugin sometimes did not update the generated code, so you had force-delete the files on every build. And then other stuff in the caches broke and you had to delete `.gradle` directory and sometimes even the `~/.gradle` directory. And sometimes the Gralde daemon hangs so you have to force it to stop with `--stop.
Go build, deno and bun are so much more reliable and faster. Something that was surprisingly fast was using the Gradle setup with skaffolding. Java hot code swapping is very fast.
> gradle is faster as a build tool for incremental compilation
Implicit:
> …than it is building from scratch, where it needs to download lots of stuff
I mean, yes, saying Java builds are fast does seem a bit “rolls eyes, yes technically by loc when the compiler is actually running” …but, ^_^! let’s not start banging on about how great cargo/rust compile times are… they’re really terrible once procedural macros are used, or sys dependencies invoke some heinous c dependency build like autoconf and lots of crates do… and there’s still a subpar incremental compilation story.
So, you know. Eh. Live and let live. Gradle isn’t that bad.
As for precedural macros - yes they can be slow to compile but so are Java annotation processors.
In my experience, java annotations gets processed very fast.
300K LOC Java should be done within ~1.5 minutes flat from zero. Can be even faster if you are running maven daemon for example. Or even within ~30 seconds if everything is within one module and javac is only invoked once.
if you do that, you should also include programming time, and divide both by the number of runs the code will have over its lifetime. Also add an appropriate fraction of the time needed to learn to program.
In such a challenge, that likely would make a very naive version win.
Apart from being impractical, I think that would be against the idea behind this challenge.
Not if you give a realistic number for an actual database of this scale. Tens of hours of dev time isn't much after you divide it by a few million.
“I discard cache and it is so slow, aarghhh”
> No external dependencies may be used
But the real meaning is that you're not allowed to use external libraries, rather than build tool related dependency.
I’ve also seen versions in Rust, Go, Python, Clickhouse and DuckDB. The discussions tab on the GitHub repo lists some of these.
Shouldn’t he take the single fastest time, assuming file (and JDK) being in file cache is controlled for?
The range of remaining three is what you expect to get 99% of the time on a real world system.
The test is very well designed, we may say.
And here is Andrei Alexandrscu arguing the same in 2012: https://forum.dlang.org/thread/mailman.73.1347916419.5162.di...
Technically yes, but these days most of my machines are single purpose VMs; database/load balancer/app server/etc, so it still seems weird not to take the fastest.
As long as you share the metal with other stuff (be it containers, binaries, VMs), there's always competition for resources, and your average time becomes your real world time.
Moreover, even if your files in memory, you cannot reserve "memory controller bandwidth". A VM using tons of memory bandwidth or a couple of cores in the same NUMA node with your VM will inevitably cause some traffic and slow you down.
There's logrotate and other cleanup tasks, monitoring, dns, a firewall, and many more stuff running on that server. No matter how much you offload to the host (or forego), there's always a kernel and supporting deamons running alongside or under your app.
The author does explicitly want all the files to be cached in memory for the later runs: https://x.com/gunnarmorling/status/1742181941409882151?s=20
Should run about 10-20% faster than that on the mentioned Hetzner hardware.
- Since we only do one decimal of floating point precision it uses integer math right from the get-go.
- FNV1-a hash with linear probing and a load factor well under 0.5.
- Data file is mmap’d into memory.
- Data is processed in 8 totally separate chunks (no concurrent data structures) and then those aggregations are in turn aggregated when all threads have finished.
I think the issues with hashing could be easily covered by having the 500 city names in the test data also be randomly generated at test time. There is no way to ensure there aren't hash collisions without doing a complete comparison between the names.
It becomes very possible to find an instance of each unique value, then runtime-design a hash algorithm where those 500 values don't collide.
Java allows self modifying code after all (and this challenge also allows native code, which can also be compiled-at-runtime)
Java is explicit, yes. This is intentional :)
I find Python to be the most pleasant personally.
For example, for performance monitoring/debugging, it's really nice to see memory usage and thread creation by package/class. Whereas, I've seen teams spend weeks trying to find memory leaks in NodeJS applications, because you just get "here's general usage by object-shape". Would literally take me less than a minute with Java...
Anyway I’m not a software engineer so take what I say about software with a grain of salt. Also I feel like other java programmers will hate you if they need to read or use your code. For example lets say you had a 2D point data type that is a record. instead of things being like “pos.add(5)” if pos is a record you need to reassign it like “pos = pos.add(5)” where add returns the pos as an object. Similar to how BigInteger works in java.
Anyway I love records because there’s zero boilerplate, it feels very similar to me to how structs are written in rust and go, or how classes are done in scala. I just never see anyone talk about it. I guess that might be because most people programming new projects on JVM are probably not programming in java?
Would have been more interesting with something like median/k-th percentile, or some other aggregation not as easy.
Or at the very least, convert the input into a more convenient binary format for the following runs.
This means that future invocations don't have to load the file from disk. This also makes it pretty critical that your program doesn't use more than 16gb of ram itself (out of the server's 32gb) or it'll push the file out of cache making future invocations of your program slower.
- Read with O_DIRECT (but compare it with the mmap approach). I know O_DIRECT gets much hate, but this could be one of the rare cases where it helps.
- Use a simple array with sentinel and linear search for the station names. This is dumb, but if the number of stations is small enough this could beat the hash. (In the back of my head I have a rule of thumb that linear search is faster for up to 1000 elements, but I'm not sure anymore where I got this from).
apparently the challenge explicitly allow for warming up the os cache on a prior run, so O_DIRECT would be disadvantaged here as the working set fits in memory.
edit: also in my very simple tests a good hashmap can beat linear search already at <10 items. It depends a lot on the hash function cost.
The sources I found corroborate this. My 1000 elements rule of thumb is apparently way off.
When you want real speed like this and are happy with insane complexity, code that writes and self-compiles a hash function based on the data can make sense.
https://gist.github.com/corlinp/176a97c58099bca36bcd5679e68f...
https://github.com/gunnarmorling/1brc#results
The fastest is currently at 12 seconds (not on M3 Pro though).
`bin/analyze measurements.txt 30.34s user 16.28s system 629% cpu 7.406 total`
With SIMD and certain assumptions about the input this can seemingly be further reduced to well under a second, eg see https://github.com/gunnarmorling/1brc/discussions/138.
I would love it if you could run my solution and compare!
https://stackoverflow.com/questions/277309/java-floating-poi...
11 seconds seems pretty impressive for a 12Gb file. Would be interesting to know what programming language could do it faster. For a database comparison you’d probably want to include loading the data into your database for a fair comparrison.
SQL in theory makes this trivial, handles many of the big optimizations and looking at the repo, cuts 1000+ LOC down to a handful. Modern SQL engines handle everything for you, which is the whole damned point of SQL. Any decent engine will handle parallelism, caching, I/O, etc. Some exotic engines can leverage GPU but for N=1 billion, I doubt GPU will be faster. Here's the basic query:
SELECT city, MIN(temp), AVG(temp), MAX(temp) FROM temps GROUP BY 1 ORDER BY 1;
In practice, generic SQL engines like PostgreSQL bloat the storage which is a Big Problem for queries like this - in my test, even with INT2 normalization (see below), pgsql took 37 bytes per record which is insane (23+ bytes of overhead to support transactions: https://www.postgresql.org/docs/current/storage-page-layout....). The big trick is to use PostgreSQL arrays to store the data by city, which removes this overhead and reduces the table size from 34GB (doesn't fit in memory) to 2GB (which does).The first optimization is to observe that the cardinality of cities is small and can be normalized into integers (INTEGER aka INT4), and that the temps can as well (1 decimal of precision). Using SMALLINT (aka INT2) is probably not faster on modern CPUs but should use less RAM, which is better for both caching on smaller systems and cache hitrate on all systems. NUMERIC generally isn't faster or tighter on most engines.
To see the query plan, use EXPLAIN:
postgres=# explain SELECT city, MIN(temp), AVG(temp), MAX(temp) FROM temps_int2 GROUP BY 1 ORDER BY 1 limit 5;
QUERY PLAN
---------------------------------------------------------------------------------------------------------------
Limit (cost=13828444.05..13828445.37 rows=5 width=38)
-> Finalize GroupAggregate (cost=13828444.05..13828497.22 rows=200 width=38)
Group Key: city
-> Gather Merge (cost=13828444.05..13828490.72 rows=400 width=38)
Workers Planned: 2
-> Sort (cost=13827444.02..13827444.52 rows=200 width=38)
Sort Key: city
-> Partial HashAggregate (cost=13827434.38..13827436.38 rows=200 width=38)
Group Key: city
-> Parallel Seq Scan on temps_int2 (cost=0.00..9126106.69 rows=470132769 width=4)
JIT:
Functions: 8
Options: Inlining true, Optimization true, Expressions true, Deforming true
(13 rows)Sigh, pg16 is still pretty conservative about parallelism, so let's crank it up.
SET max_parallel_workers=16; set max_parallel_workers_per_gather=16;
SET min_parallel_table_scan_size=0; set min_parallel_index_scan_size=0;
SET parallel_setup_cost = 0; -- Reduce the cost threshold for parallel execution
SET parallel_tuple_cost = 0.001; -- Lower the cost per tuple for parallel execution
postgres=# explain SELECT city, MIN(temp), AVG(temp), MAX(temp) FROM temps_int2 GROUP BY 1 ORDER BY 1 limit 5;
...
Workers Planned: 14
...
top(1) is showing that we're burying the CPU: top - 10:09:48 up 22 min, 3 users, load average: 5.36, 1.87, 0.95
Tasks: 169 total, 1 running, 168 sleeping, 0 stopped, 0 zombie
%Cpu(s): 12.6 us, 4.8 sy, 0.0 ni, 7.8 id, 74.2 wa, 0.0 hi, 0.7 si, 0.0 st
MiB Mem : 32084.9 total, 258.7 free, 576.7 used, 31249.5 buff/cache
MiB Swap: 0.0 total, 0.0 free, 0.0 used. 30902.4 avail Mem
PID USER PR NI VIRT RES SHR S %CPU %MEM TIME+ COMMAND
1062 postgres 20 0 357200 227952 197700 D 17.9 0.7 9:35.88 postgres: 16/main: postgres postgres [local] SELECT
1516 postgres 20 0 355384 91232 62384 D 17.6 0.3 0:08.79 postgres: 16/main: parallel worker for PID 1062
1522 postgres 20 0 355384 93336 64544 D 17.6 0.3 0:08.53 postgres: 16/main: parallel worker for PID 1062
1518 postgres 20 0 355384 90148 61300 D 17.3 0.3 0:08.53 postgres: 16/main: parallel worker for PID 1062
1521 postgres 20 0 355384 92624 63776 D 17.3 0.3 0:08.54 postgres: 16/main: parallel worker for PID 1062
1519 postgres 20 0 355384 90440 61592 D 16.6 0.3 0:08.58 postgres: 16/main: parallel worker for PID 1062
1520 postgres 20 0 355384 92732 63884 D 16.6 0.3 0:08.49 postgres: 16/main: parallel worker for PID 1062
1517 postgres 20 0 355384 91544 62696 D 16.3 0.3 0:08.55 postgres: 16/main: parallel worker for PID 1062
interestingly, when we match workers to CPU cores, we don't %CPU drops to 14% i.e. we don't leverage the hardware.OK enough pre-optimization, here's the baseline:
postgres=# SELECT city, MIN(temp), AVG(temp), MAX(temp) FROM temps_int2 GROUP BY 1 ORDER BY 1 limit 5;
city | min | avg | max
------+-----+----------------------+------
0 | 0 | 276.0853550961625011 | 1099
1 | 0 | 275.3679265859715333 | 1098
2 | 0 | 274.6485567539599619 | 1098
3 | 0 | 274.9825584419741823 | 1099
4 | 0 | 275.0633718875598229 | 1097
(5 rows)
Time: 140642.641 ms (02:20.643)
I also tried to leverage a B-tree covering index (CREATE INDEX temps_by_city ON temps_int2 (city) INCLUDE (temp) )
but it wasn't faster - I killed the job after 4 minutes. Yes, I checked that it pg16 used a parallel index-only scan
(SET random_page_cost =0.0001; set min_parallel_index_scan_size=0; set enable_seqscan = false;
SET enable_parallel_index_scan = ON;) - top(1) shows ~1.7% CPU, suggesting that we were I/O bound.Instead, to amortize the tuple overhead, we can store the data as arrays:
CREATE TABLE temps_by_city AS SELECT city, array_agg(temp) from temps_int2 group by city;
$ ./table_sizes.sh
...total... | 36 GB
temps_int2 | 34 GB
temps_by_city | 1980 MB
Yay, it now fits in RAM. -- https://stackoverflow.com/a/18964261/430938 adding IMMUTABLE PARALLEL SAFE
CREATE OR REPLACE FUNCTION array_min(_data ANYARRAY) RETURNS NUMERIC AS $$
SELECT min(a) FROM UNNEST(_data) AS a
$$ LANGUAGE SQL IMMUTABLE PARALLEL SAFE;
SET max_parallel_workers=16; set max_parallel_workers_per_gather=16;
SET min_parallel_table_scan_size=0; set min_parallel_index_scan_size=0;
SET parallel_setup_cost = 0; -- Reduce the cost threshold for parallel execution
SET parallel_tuple_cost = 0.001; -- Lower the cost per tuple for parallel execution
postgres=# create table tmp1 as select city, min(array_min(array_agg)), avg(array_avg(array_agg)), max(array_max(array_agg)) from temps_by_city group by 1 order by 1 ;
SELECT 1001
Time: 132616.944 ms (02:12.617)
postgres=# explain select city, min(array_min(array_agg)), avg(array_avg(array_agg)), max(array_max(array_agg)) from temps_by_city group by 1 order by 1 ;
QUERY PLAN
-------------------------------------------------------------------------------------------------
Sort (cost=338.31..338.81 rows=200 width=98)
Sort Key: city
-> Finalize HashAggregate (cost=330.14..330.66 rows=200 width=98)
Group Key: city
-> Gather (cost=321.52..322.64 rows=600 width=98)
Workers Planned: 3
-> Partial HashAggregate (cost=321.52..322.04 rows=200 width=98)
Group Key: city
-> Parallel Seq Scan on temps_by_city (cost=0.00..0.04 rows=423 width=34)
(9 rows)Ah, only using 3 cores of the 8...
---
CREATE TABLE temps_by_city3 AS SELECT city, temp % 1000, array_agg(temp) from temps_int2 group by 1,2;
postgres=# explain select city, min(array_min(array_agg)), avg(array_avg(array_agg)), max(array_max(array_agg)) from temps_by_city3 group by 1 order by 1 ;
QUERY PLAN
---------------------------------------------------------------------------------------------------------
Finalize GroupAggregate (cost=854300.99..854378.73 rows=200 width=98)
Group Key: city
-> Gather Merge (cost=854300.99..854348.73 rows=2200 width=98)
Workers Planned: 11
-> Sort (cost=854300.77..854301.27 rows=200 width=98)
Sort Key: city
-> Partial HashAggregate (cost=854290.63..854293.13 rows=200 width=98)
Group Key: city
-> Parallel Seq Scan on temps_by_city3 (cost=0.00..98836.19 rows=994019 width=34)
JIT:
Functions: 7
Options: Inlining true, Optimization true, Expressions true, Deforming true
(12 rows)This 100% saturates the CPU and runs in ~65 secs (about 2x faster).
---
Creating the data
Here's a very fast lousy first pass for PostgreSQL that works on most versions - for pg16, there's random_normal(). I'm not on Hetzner so I used GCP (c2d-standard-8 with 16vCPU, 64GB, and 150GB disk, ubuntu 22.04 and
CREATE TABLE temps_int2 (city int2, temp int2);
-- random()*random() is a cheap distribution. 1100 = 110 * 10 to provide one decimal of precision.
INSERT INTO temps_int2 (city, temp) SELECT (1000*random())::int2 as city, (random()*random()*1100)::int2 temp from generate_series(1,1e9)i;1) use the integer version of generate_series, the 1e9 leads to the floating point version being chosen
2) move the generate_series() to the select list of a subselect - for boring reasons, that we should fix, the FROM version materializes the result first
3) using COPY is much faster, however a bit awkward to write
psql -Xq -c 'COPY (SELECT (1000random())::int2 as city, (random()random()*1100)::int2 temp FROM (SELECT generate_series(1,1e9::int8))) TO STDOUT WITH BINARY' | psql -Xq -c 'COPY temps_int2 FROM STDIN WITH BINARY'
1/10th scale tests:
psql -c 'create table temps_int2_copy as select * from temps_int2_copy where 1=0;'
psql -Xq -c 'COPY (SELECT (1000*random())::int2 as city, (random()*random()*1100)::int2 temp FROM (SELECT generate_series(1,1e8::int8))) TO STDOUT WITH BINARY' | psql -Xq -c 'COPY temps_int2_copy FROM STDIN WITH BINARY'
==> 32sec psql -c 'create table temps_int2_copy2 as select * from temps_int2_copy where 1=0;'
psql -c 'INSERT INTO temps_int2_copy2 (city, temp) SELECT (1000*random())::int2 as city, (random()*random()*1100)::int2 temp from generate_series(1,1e8::int)i;'
==> 90sec psql -c 'create UNLOGGED table temps_int2_copy3 as select * from temps_int2_copy where 1=0;'
psql -c 'INSERT INTO temps_int2_copy3 (city, temp) SELECT (1000*random())::int2 as city, (random()*random()*1100)::int2 temp from generate_series(1,1e8::int)i;'; date
==> 45sec (still not faster!)Of course, if we really want to "go fast" then we want parallel loading, which means firing up N postgresql backends and each generate and write the data concurrently to different tables, then each compute a partial summary in N summary tables, and finally merge them together.
echo "COPY temps_int2_copy_p1 from '/var/lib/postgresql/output100m.txt' with csv;" | /usr/lib/postgresql/16/bin/postgres --single -D /etc/postgresql/16/main/ postgres
==> 32sec (saturates one core)The ultimate would be to hack into postgres and skip everything and just write the actual filesystem files in-place using knowledge of the file formats, then "wire in" these files to the database system tables. This normally gets hairy (e.g. TOAST) but with a simple table like this, it might be possible. This is a project I've always wanted to try.
We (postgres) should fix that at some point... The difference basically is that there's a dedicated path to insert many tuples at once that's often used by COPY that isn't used by INSERT INTO ... SELECT. The logic for determining when that optimization is correct (consider e.g. after-insert per-row triggers, the trigger invocation for row N may not yet see row N+1) is specific to COPY right now. We need to generalize it to be usable in more places.
To be fair, part of the reason the COPY approach is faster is that the generate_series() query actually uses a fair bit of CPU on its own, and the piped psql's lead to the data generation and data loading being run separately. Of course, partially paying for that by needing to serialize/deserialize the data and handling all the data in four processes.
When doing the COPYs separately to/from a file, it actually takes longer to generate the data than loading the data into an unlogged table.
# COPY (SELECT (1000*random())::int2 as city, (random()*random()*1100)::int2 temp FROM (SELECT generate_series(1,1e8::int8))) TO '/tmp/data.pgcopy' WITH BINARY;
COPY 100000000
Time: 21560.956 ms (00:21.561)
# BEGIN;DROP TABLE IF EXISTS temps_int2; CREATE UNLOGGED TABLE temps_int2 (city int2 NOT NULL, temp int2 NOT NULL); COPY temps_int2 FROM '/tmp/data.pgcopy' WITH BINARY;COMMIT;
BEGIN
Time: 0.128 ms
DROP TABLE
Time: 0.752 ms
CREATE TABLE
Time: 0.609 ms
COPY 100000000
Time: 18874.010 ms (00:18.874)
COMMIT
Time: 229.650 ms
Loading into a logged table is a bit slower, at 20250.835 ms.> Of course, if we really want to "go fast" then we want parallel loading, which means firing up N postgresql backends and each generate and write the data concurrently to different tables, then each compute a partial summary in N summary tables, and finally merge them together.
With PG >= 16, you need a fair bit of concurrency to hit bottlenecks due to multiple backends loading data into the same table with COPY. On my ~4 year old workstation I reach over 3GB/s, with a bit more work we can get higher. Before that the limit was a lot lower.
If I use large enough shared buffers so that IO does not become a bottleneck, I can load the 1e9 rows fairly quickly in parallel, using pgbench:
c=20; psql -Xq -c "COPY (SELECT (1000*random())::int2 as city, (random()*random()*1100)::int2 temp FROM (SELECT generate_series(1,1e6::int8))) TO '/tmp/data-1e6.pgcopy' WITH BINARY;" -c 'DROP TABLE IF EXISTS temps_int2; CREATE UNLOGGED TABLE temps_int2 (city int2 NOT NULL, temp int2 NOT NULL); ' && time pgbench -c$c -j$c -n -f <( echo "COPY temps_int2 FROM '/tmp/data-1e6.pgcopy' WITH BINARY;" ) -t $((1000/${c})) -P1
real 0m26.486s
That's just 1.2GB/s, because the bottleneck is the per-row and per-field processing, due to their narrowness.> The ultimate would be to hack into postgres and skip everything and just write the actual filesystem files in-place using knowledge of the file formats, then "wire in" these files to the database system tables. This normally gets hairy (e.g. TOAST) but with a simple table like this, it might be possible. This is a project I've always wanted to try.
I doubt that will ever be a good idea. For one, the row metadata contain transactional information, that'd be hard to create correctly outside of postgres. It'd also be too easy to cause issues with corrupted data.
However, there's a lot we could do to speed up data loading performance further. The parsing that COPY does, uhm, show signs of iterative development over decades. Absurdly enough, that's where the bottleneck most commonly is right now. I'm reasonably confident that there's at least 3-4x possible without going to particularly extreme lengths. I think there's also at least a not-too-hard 2x for the portion of loading loading data into the table.
I think it actually shows that you're IO bound (the 'D' in the 'S' column). On my workstation the query takes ~10.9s after restarting postgres and dropping the os caches. And this is a four year old CPU that wasn't top of the line at the time either.
> postgres=# explain select city, min(array_min(array_agg)), avg(array_avg(array_agg)), max(array_max(array_agg)) from temps_by_city group by 1 order by 1 ;
Note that you dropped the limit 5 here. This causes the query to be a good bit slower, the expensive part here is all the array unnesting, which only needs to happen for the actually selected cities.
On my workstation the above takes 1.8s after adding the limit 5.
It does, but I restarted postgres and cleared the OS cache.
Btw, the primary bottleneck for the array-ified query is the unnest() handling in the functions. The minimal thing would be to make the functions faster, e.g. via:
CREATE OR REPLACE FUNCTION array_avg(_data anyarray) RETURNS numeric IMMUTABLE PARALLEL SAFE LANGUAGE sql AS $$ select avg(a) from (SELECT unnest(_data) as a) $$;
But that way the unnest is still done 3x. Something like
SELECT a_min, a_max, a_avg FROM temps_by_city, LATERAL (SELECT min(u) a_min, max(u) a_max, avg(u) a_avg FROM (SELECT unnest(array_agg) u)) limit 5;
should be faster.852 ms on an M1 Max
This is for use on a BSD where Java may not work too well.
Thanks
It’s just the city names + averages from the official repository using a normal distribution to generate 1B random rows.
https://nestedsoftware.com/2018/03/20/calculating-a-moving-a...
so you just walk the file and read a chunk, update the averages and move on. the resource usage should be 0.000 nothing and speed should be limited by your disk IO.
So it's not IO-bound.
BTW asking ChatGPT to utilize all cores did not yield anything working in reasonable time.
time duckdb -list -c "select map_from_entries(list((name,x))) as result from (select name, printf('%.1f/%.1f/%.1f',min(value), mean(value),max(value)) as x from read_csv('measurements.txt', delim=';', columns={'name': 'varchar', 'value':'float'}) group by name order by name)"
takes about 20 seconds