How We Built Uber’s Highest Query per Second Service Using Go
eng.uber.com
eng.uber.com
More than 3 years ago we've implemented a reverse geocoder web-service that indexed complete Census TIGER dataset. The service handled over 15K reqs/sec on a 2011 macbook, doing exact in polygon search(no approximations). We implemented an optimized (for in polygon search) R-tree data structure for the lookups.
Assuming Uber's geofence lookups would rely on a much smaller dataset than TIGER, I think you could have come up with a much more efficient implementation requiring a fraction of the resources.
And again, use of Go; irrelevant.
This is a scientific paper.
I wouldn't be surprised that due to the lack of generics in golang, this approach was less feasible.
Also, how does it make testing harder?
You can't typecheck the generic itself. Especially conversions could be really hard with them.
(I love generics, still this somewhat aweful). Currently Converting Generics between Scala and Java is totally Ugly.
(disclaimer: i think that lack of generics in Go is a mistake)
Go generate, for example, brings me back memories of a time C++ compilers used pre-processor magic for generic programming, around 1994 or so.
Something like
#define Type1 my nice type1 definition
#include <generic-list.h>
#undef Type1
#define Type1 my nice type2 definition
#include <generic-list.h>
ListType1 ...
ListType2 ...This makes sense, and when I first started writing Go, I felt the same way, but almost four years later, concurrency is one of the least important reasons I still write Go.
I wouldn't say all Gophers feel this way, but I know it's a very common experience for Gophers - "you come for the concurrency, but stay for the interfaces"[0].
[0] I've alternatively heard things like "readability", "tooling", and "robustness" used in the place of "interfaces" in this quote.
Yup. I literally just gave a talk at two Go conferences in the last two weeks (GopherCon India and GopherCon Dubai) about the io.Reader/io.Writer interfaces specifically. They're deceptively simple on the surface, but they're insanely powerful once you 'get' them.
There was a blog post a while ago about how golangs's interfaces caused issues in production because they're implicit.
Check out type classes for a superior way to solve this issue (e.g. what Scala or Rust do).
They have an N readers, one writer problem. It's possible to do that without blocking the readers. There was a YC article on a lockless solution to that yesterday. This is an easier problem. Each city's geofences can be updated by making a copy, updating the copy, and atomically replacing the single reference to the copy. Any searches still running on the old copy continue to run on the old copy. Eventually, the GC deletes the old copy.
Why? Concurrent data structures are extremely useful.
> The code is more readable and you don't offload any work to the GC.
You shouldn't write concurrent data structures yourself; you should use a library. And the amount of time spent in the GC for this is going to be negligible assuming writes are infrequent compared to reads. It's rare in a GC'd language that your write operation won't be creating garbage somewhere along the way, so it ends up being amortized.
Which one are you referring to?
"While the runtime complexity of the solution remains O(N), this simple technique reduced N from the order of 10,000s to the order of 100s."
It seems odd to me that they're posting about the performance of Go, yet they deliberately chose a less optimal algorithm because the better ones were "complicated". Or, if their simpler approach did in fact run faster than R-trees or S2, it would have been nice to see some benchmarks and a clearer explanation for why. Choosing the best algorithm seems more important than the language in a case like this one.
I always have fun trotting out Radix Sort and watching fresh CS grads try to understand how a O(kN) walks over O(NlogN) by a factor of 10x+.
Fun bit about Radix is it really scales well if your dataset fits in memory. Last time I benched my lame Java implementation it was about 2-5x faster than the built in Arrays.sort() on native values(float, int) and 20x faster when you started putting it up against Comparable. Trended that way well up to 65k+ entries.
It gets even better if you can drop down to a proper native language that lets you prefetch.
Using O analysis is for seeing how an algorithm scales with respect to the input size. If anything, the fact that someone believes a O(nlogn) algorithm is always faster than O(n^2) is a failure on that person's part in properly understanding what O is used to measure.
If you can come up with a goal of what you need the response times to be, I think the easiest (in terms of readability, testability, length, etc) solution to meet the goal should win.
That being said sometimes brute force is good enough.
Also don't know why they chose raycasting over winding number.
I've been working on geo stuff in golang and building the geometry libraries for spatial indexing was easy enough to do in my spare time, so I don't know why they couldn't at least try.
Source: just spent two months implementing R-trees in Rust
Geo data is generally nicely distributed. R-tree (or similar) is the correct solution here.
Or the fact that due to the lack of generics in golang, it would make things even more complicated, and perhaps losing performance as well.
We make pretty extensive use of R-Tree indexes in PostGIS, tables with 10,000,000+ rows and shapes that roughly match streets, parcels, blocks, zip codes. While I don't have 99th percentile data offhand average query time is <1ms. Perhaps there are slow responses at the 99th percentile however I would be very surprised by this.
EG: Subdivide a flat projection of earth into n^2 squares. Create an array of length n^2. Set the value of each element in the array to a list of canidate geofences(which have area in that square).
Scale lat and long between 0 and 1. Then you can index directly into it with PrecomputedArray[floor(lat*n+long)]
This is trading space for time, may as well choose space here.
Great approach!
It feels like this type of narrative is becoming more common.
One creates the first version in Node.js/Python as its quick to build and iterate an app. At this stage you are still hammering out exactly how things need to be done and accessing the actual needs of your users. Speed of iteration is the most important attribute for any choice made during this stage as generally you don't have enough scale for choices to matter much and it prevents wasted effort on things that end up on the cutting room floor.
Once you have built the app, and it gets stable major feature wise, you then have plenty of data to drive decisions about what the pieces actually need to be and the tool best suited for them. Go happens to be well suited for swapping out various generic back end pieces with something custom to the problem at hand.
Another way to think about it is the Node.js/Python version is the next step up from diagramming the app out on paper. Its more like the previsulazation stage of a movie (roughly animated, with intern voiceovers). This is important as it can save you a bunch of time on things that are unneeded and can help identify problem areas you should focus on first.
That's the narrative, but I don't find Python (I have minimal experience with Node) faster to write in any way but library availability, and very painful to refactor without the ability to reach for a type checker.
My question isn't silly. I just now wrote two short programs, one in Go and one Javascript, to run a vector multiply and add on two length 1000 vectors, a million times. The JS program is faster.
Result (in seconds):
$ node vmadd.js
0.984
$ go run vmadd.go
1.141044662
JS code:
https://gist.github.com/hwinkler/62c9da0d9f2d8c7981b0Go code: https://gist.github.com/hwinkler/d234bb62b2c5fa081a50
$ go build vmadd.go
$ ./vmadd
1.106074715
This is a 2013 Core i7 15" MBP'Highest QPS?'
This is a trivial geometric problem, any sane implementation should be several orders of magnitude faster than doing network IO. /rant
"Instead of indexing the geofences using R-tree or the complicated S2..."
"... we first find the desired city with a linear scan of all the city geofences"
Why not use a spatial index? It's not hard, and you wouldn't need to worry about rebuilding your index because your city geofences are not likely to change frequently.
The bottleneck here isn't rooted in I/O but a better algorithmic approach to the problem.
Can't you precompute the X and Y minima and maxima for each geo fence and throw out 99% of candidates extremely quickly? In other words store the bounding rectangle for each geofence. Even with 100,000s of entries we're talking 4 primitive type comparisons per entry to exclude all non-possible candidates. I would have a hard time believing this is slower than ray casting to city geofences.
Check 100 000 entries (4 comparisons each) - > filter out 1% to do expensive calculations on.
Vs
Check 25 cities (4 comparisons each) - > filter out 1/25 Do another easy run on the filtered 4000 and get the same areas to do expensive calculations on as above.
That's 25+4000 easy calculations against 100 000.
But yeah, I'm not sure why they seem to go the hard route from the beginning.
If you have a nested series of rectangular bounding boxes you're well on your way to building an R-tree.
At least the article could have addressed why that would not work for their use case.
edit: also now they have to classify polygons into groups that wholly fit into some "city" polygon. That seems like it becomes a pain to maintain since you may want to expand a region in the future and it happens to go outside of your previously defined "city". Or maybe you want a region that isn't in a "city", like a "country" or "mid-atlantic" region.
Eg. latency across availability zones in the same EC2 region is <2ms (and obviously even faster in the same AZ)
That's misleading, math in V8 can be VERY fast (aside from the current absence of vectorization). See the benchmarks here: http://julialang.org
I suspect the lack of an integer primitive is a big deal in both Lua and JavaScript.
Not as much as you'd think. All JavaScript JITs speculate that numbers that have been observed to be integers remain integers and compiles them to use integer operations. As long as your code is hot enough to enter the optimizing JIT, the overhead is solely in whatever overflow checks couldn't be optimized out.
http://benchmarksgame.alioth.debian.org/u64q/performance.php...
http://benchmarksgame.alioth.debian.org/u64q/performance.php...
However, a mutex lock/unlock should take <100ns. It doesn't seem like that should be a bottleneck. It'd be interesting to hear more about the data, requirements, and response times (median, 90%, etc).
Can't remember the exact results but throughput was much higher than what I was able to achieve using highly optimized PostGIS queries.
Yes, it has. Post-1.0 Rust's design philosophy is about taming shared state concurrency: that is, disallowing data races statically, making sure you take locks or use atomics where possible, and having a robust set of generic concurrent data structures ready for use so you don't have to implement your own.
This article, for example, rightly observes that atomics are difficult to get right, making it not scalable to have every program implement concurrent data structures from scratch. The benefit of having a rich set of generic concurrent data structures readily available is that it mitigates this problem.
In retrospect it was clear that channels were very elegant for certain parts of the code, and in other areas mutexes would have been simpler.
After a bit of experience with go, you gain an intuitive sense for which is appropriate in each case.
Channels are good in other scenarios.
Similar to other comments, the article is a bit lacking in detail. Regardless of language used, I would have liked to see a comparison of query times using a spatial index and not. From my limited experience with PostGIS, there are some inside polygons that perform poorly on large polygons, where Mongodb has done better, and vice versa with linestrings.
It's ~4k per instance. I would like to know what the R/W ratio is. And how often the index is updated is the key, aka locking operations, which should be mentioned in the article.
Also this article should give more comparison on other implementations, if possible. At least compare with Uber's monolithic implementation.
The geo-datastructure part is very cool.
cluster's been around a while now, and is built in to node 5.x.
I can buy that the language is easy to learn since the feature set is very limited, but is productivity in the first week/month really what to optimise for?
They also list performance as a benefit.
Go is a reasonable choice, but I doubt that the result would have been much different if they'd used any other modern statically typed language such as Rust, Scala, D, or even Nim.
I'm not the biggest node fan, but that is utterly false. How do you think asynchronous system calls get processed in libuv (the event loop that node uses.) ? threads. Of course you'll park the background thread while you're waiting in case of IO, but if it's CPU intensive then you bet that you can use all cores with background work. You just need to know how to write C++ and read the documentation on the bindings
Saying "of course node.js supports multiple threads of concurrent computation, you just have to write in c++ instead of javascript" is missing the point of why people use it.