Handling 1M Requests per Minute with Go
marcio.io
marcio.io
1) Take the original code, do the upload exactly in place in the original request (not even spawning a goroutine). However: protect the upload with a semaphore which only allows N-in-flight.
My reasoning is, well, if the system operates with low latency when operating nominally, blocking the incoming request isn't too painful. The reason there was a problem in the first place that there were too many requests in flight and the system hit a meta-stable state where no requests could complete efficiently.
2) (or instead of (1)): If you're going to have a worker pool, why have that complicated chan-chan-Job business? It seems that `func StartProcessor` was close to being a viable solution. All you need is to start a few of those in parallel, each reading from the same `Queue`. Was there a reason to introduce the `WorkerPool chan chan Job`? That looks quite a bit more complicated than it needs to be. The queues don't need to be separate per worker unless there is some other substantial reason.
--
The next thing one would need to take care of is to ensure that the whole system doesn't stall due to a broken/laggy network, so, to put some timeouts on the S3 uploads, for example, to ensure the system can return to a stable state on its own when the thundering herd has passed.
Re: Queue: Wouldn't the queue involve locking lest two workers end up trying to work on the same request? To be completely concurrent, I guess one could use a lock-free data structure instead (or implement one on top of something like RocksDB)?
[0] http://ferd.ca/queues-don-t-fix-overload.html
[0] http://engineering.voxer.com/2013/09/16/backpressure-in-node...
Block on N-Semaphore, with timeout
Do timeout upload
Replace N-Semaphore
If you don't always replace the semaphore, that's a bug.Queue: I'm just comparing to what the article does. It already has contention on a queue (the chan), it's just the chan-chan-Worker rather than chan-Job. In practice, go channels happily handle millions of messages contending to multiple workers just fine. Consider this test example where you aren't even actually burning any CPU to perform the work:
package main
func main() {
q := make(chan int)
for i := 0; i < 10; i++ {
go func() {
for x := range q {
x = x * 10
}
}()
}
for i := 0; i < 1000000; i++ {
q <- i
}
}
On my laptop, it runs in 0.333s single core, and it's slightly slower when you set GOMAXPROCS > 1. But not much slower, the total runtime goes to 0.4-0.5s or so. (Measured with go 1.4). As soon as you do any actual work with the messages you are passing around, the overhead of locking will be lost in the noise.Don't get me wrong, the most interesting part is definitely the implementation and as a Go noob I found it very useful - it's just a bit misleading for the headline to sum your request rate across all parallelized machines.
I don't get this reasoning.
I agree that other languages could accomplish the same effect, but I'm not surprised by their sentiment given where they were coming from.
It makes perfect sense to me. What would you have recommended them, for a reasonably-high-performance server implementation? (Please don't say C.)
https://benchmarksgame.alioth.debian.org/u64q/benchmark.php?...
Yes. Other than some very basic scripts, I have never worked with Ruby, so I was not aware that it is so slow. Thanks for pointing that out.
As for my recommendation, it is pretty standard worker architecture:
Their system seems to be an ingesting-only system, that is, the clients are getting an empty HTTP 200 OK response. Given this, I would put openresty (nginx) in the front, with some trivial Lua code[1] to en-queue payloads to beanstalkd. Then, you can either have your workers inside openresty (using Openresty timers) or have them as separate processes and written in the language of choice. We have been using this for a couple of years now and it is working really well for our use case, also an ingesting-only system.
[1] https://github.com/smallfish/lua-resty-beanstalkd/blob/maste...
Than MRI Ruby you mean.
The advantage wouldn't be as much if Ruby designers cared to add AOT compilation in the same vein as Dylan or Common Lisp to the canonical implementation.
> Than MRI Ruby you mean.
Given the nature of orders-of-magnitude comparisons and the lack of Ruby implementations that are even one order of magnitude faster than MRI, "...than Ruby" is reasonably accurate if "...than MRI Ruby" is at all accurate.
> The advantage wouldn't be as much if Ruby designers cared to add AOT compilation in the same vein as Dylan or Common Lisp to the canonical implementation.
Maybe, though that's unproven. AFAIK, actual Ruby implementations with AOT only seem to gain about a factor of 2 improvement, not an order of magnitude.
You could have written this in any number of languages. I don't know why go is more logical than say Java JavaScript.
That may be true for some types of code, but for an app that's predominantly shuffling data over the network you should be spending most of the time in kernel space executing syscalls, and then language differences are largely irrelevant.
> and much easier to write concurrent code in.
How? Writing concurrent code in Ruby is trivial since 1.9.x (prior to 1.9 you had to battle the green threads in MRI for some stuff), which isn't exactly new.
Having spent time recently optimizing a setup that does pretty much exactly what they're doing (in Rails; I don't like Rails, but in this case Rails is not a problem), the time saved doing things like eliminating unnecessary buffering all over the place (e.g. POSTs gets buffered by pretty much everything that likes to consider itself a web server, sometimes multiple times; if you e.g. run Nginx in front of pretty much any Ruby web servers running Rails, worst case you may end up with things passed through at least 3, possibly 4 buffers).
Basically they should be IO bound.
Based on their description, they should be IO bound. If they're IO bound, their system should spend the vast majority of time in the kernel executing system calls.
If they do that, then whichever language they use has very little relevance to whether or not their system ends up being fast.
EDIT: As an example, my first ever production Ruby app was a messaging middleware server that processed about 3 million messages a day (34/sec) on 10% of a single mid-range 2005-era Xeon core. It was totally unoptimized, and used select() instead of the more modern, faster alternatives. Also, only about 10% of that time was spent in user space. 90% of it was in the kernel, executing system calls, so most meaningful optimization would be in improving the usage of system calls that are exactly the same across languages.
At the time, going beyond a single core required multiple instances, as Ruby lacked native threads back then. But processing 500-600/sec across 2-3 processes on a dual core 2005-era Xeon would not have been challenging even then with that quite naive and dated approach.
As someone else pointed out, they're doing ~4k/sec on a modern dual-core Xeon. Based on the above I'd expect it to be fairly easy to match their ~4K/sec with Ruby today.
Suggestion: Make an S3-uploader package with _internal_ connection pooling, upload queueing and concurrency handling.
Nice features include:
> Retry Everything: All http requests and every part is retried on both uploads and downloads.
> Configurable conncurrency
> Uses an io.Writer (you could actually start posting to S3 before all of the data gets in on your side.)
etc.
the new garbage collector in 1.5 should improve things.
https://talks.golang.org/2015/state-of-go-may.slide#6 (Slides 6 to 11)
When tuning a prod Scala deployment a while back, I encountered a nasty "pregnant pause" every once in a while similar to what you have experienced. So, below are what I used in annotated form and adapted for the deployment you describe (the max heap setting). Some of these may be completely obvious to you (or others reading), yet are included for completeness.
-server
This one should be obvious :-)
-Xmx54G
Memory is cheap, so give the JVM 54 gig so that GC
isn't forced to run when your system is in the
steady-state of 48G heap utilization.
-XX:PermSize=128m -XX:MaxPermSize=1024M
Ditto on the cheapness of memory.
-Xss1M
A stack size of 1 meg seems a bit much, but does
allow for recursive algorithms to operate with
impunity.
-XX:ReservedCodeCacheSize=128m
This one was needed for Scala. It likely should
be specified with a high value since it limits the
JIT's code cache.
-XX:+DoEscapeAnalysis
A nice way to releave some heap pressure[1].
-XX:+UseCodeCacheFlushing
Should the ReservedCodeCacheSize be exceeded, this
lets the JIT continue to do its thing in an LRU
type of fashion (I believe).
-XX:+UseParallelGC
This one is the most impactful one of all. A lot
of people will say "use UseConcMarkSweepGC!" They
are wrong for high volume server deployments. The
concurrent mark and sweep algorithm caused massive
"pregnant pauses" in prod for me! The Parallel GC
algorithm performs much better under load and
doesn't cause the VM to sit-and-spin for 20+
seconds.
-XX:+UseCondCardMark
Another tweak which had a major performance boost
for me[2].
-XX:+UseNUMA
If your servers are NUMA[3] based, then this can
significantly increase performance[4] as well.
HTH1 - http://www.ibm.com/developerworks/java/library/j-jtp09275/in...
2 - https://blogs.oracle.com/dave/entry/false_sharing_induced_by...
3 - https://en.wikipedia.org/wiki/Non-uniform_memory_access
4 - http://jose-manuel.me/2011/06/numa-bb/
EDIT: Inserted newlines to eliminate horizontal scrolling.
IIRC, I did and it didn't benchmark well for my needs. However, it all depends on what JRE you're using and what the system is doing to pick the GC which is best for any given deployment. Classic case of YMMV and all that.
In the end, even though it may sound trite, the only way to know what works best for a given combination of JRE/OS/hardware is to measure it. This article[1] had some good tips and Mission Control[2] is a huge help in this arena.
1 - http://www.infoq.com/articles/Tuning-Java-Servers?utm_source...
2 - http://docs.oracle.com/javacomponents/jmc-5-5/jmc-user-guide...
Tuning GC for Spark: https://www.youtube.com/watch?v=drmJDISLkf4
In the Big Data space we have dozens of machines all with very large stack sizes (I run mine with 250GB) and don't run into any major stop the world pauses.
But I only maintain EE apps on at most 10gb heaps, so I'm not hugely experienced with tuning this. All I can recommend is heap size and CMS thresholds (no G1 exp) set so that at steady state you don't end up hitting full GCs.
Are there any tutorials/templates/best practices for writing a small web service in Golang?
[0] https://golang.org/doc/articles/wiki/
[1] https://github.com/gin-gonic/gin
I loved it when I started, then really hated it when I needed to refactor, there is just too much magic and you pay a speed cost for that magic.
I do enjoy coding in Golang, but we use mostly Java where I work, and for us, the benefits don't make up for the things we lose. This blog post is a great example: the solution they had to find is the first thing you'd probably do in Java, because Java has a standard package with all sorts of concurrency patterns.
Go really needs a library for these patterns built in... I assume the lack of generics prevents users from creating that themselves (I'm not trying to start a language war here, seriously).
Except Go provides limited value for somebody who already knows .net or a jvm language. And that's a huge chunk of the market.
And the key advantage that Go offers to business is that they can easily hire from a pool of experienced programmers and have them learn Go with little downtime.
1).Initialize a job channel
2).Initialize a set of workers that listen to this channel to pull the jobs indefinitely. In this case just call go StartProcessor() for fixed number of times.
What confuses me is that IMO workPoolChannel isn't necessary here. What is the consideration behind to use a channel for workers?
Without really knowing the company's needs, I am relying on this paragraph from the post:
While working on a piece of our anonymous telemetry and analytics system, our goal was to be able to handle a large amount of POST requests from millions of endpoints. The web handler would receive a JSON document that may contain a collection of many payloads that needed to be written to Amazon S3, in order for our map-reduce systems to later operate on this data.
Knowing this, I would build it differently.
1. Clients post to S3 Directly 2. Lambda -> Overload business logic, private data, cleanup, spam control etc... 3. Prepare files (64M) for Hadoop 4. Hadoop
There's no reason to have that proxy in the middle, Amazon S3 will handle those millions of requests with no real trouble, I wouldn't throw machines on this process.
The difference is management of a simple process (behind elastic load balancer, of course) vs a more complicated architecture with three distinct, load balanced process types (webserver -> queue -> worker).
https://twitter.com/julianobs/status/614416512825323520
Hey, and Elixir is already way more expressive than Go and it's incredibly easy to build fault-tolerant and distributed systems, not to mention the productivity gains when using the phoenix framework!
Seriously, posting requests/second metric without any context about hardware and sample code doesn't help anyone.
but hey, as long as stuff responds in microseconds with zero errors under load, just use it ! (Also, clojure is a way better language than go, too.)
16k requests per second is worth writing about if you're talking about a process with substantial side effects (S3 I/O).
They're doing development without understanding how computers work, where the bottlenecks are, or what the maximum theoretical throughput for the use-case is. They ended up with something slightly better than the horrible situation they were in, and are celebrating a inefficient solution as a technical triumph.
They were able to solve their problem in a single process balanced over 4 boxes without ever having to hire someone like you, despite your expertise.
Could they have increased throughput? Absolutely. It would have involved a different architecture with more complexity & time, and it also would have relied on skills beyond what was immediately available. I'm guessing their line count is around ~200 for the core functionality.
Can you share some actual technical points where they made an error? I would really like to see you demonstrate expertise beyond these uninspiring generalities.
Well, Go's runtime allocates 8KB (last I recall) of growable stack space per goroutine. Assuming that their first solution was deployed on the same instance type as their final solution: c4.large (3.75 GB), then they could handle at most ~470,000 outstanding gorountines; assuming that all RAM is used for only gorountines, which is not realistic of course. So their server fell over once it exhausted memory.
This type of memory exhaustion isn't a problem with evented I/O. You have a single thread that responds to async events related to the I/O you're performing.
So, due to the limitations of Go's runtime, they settled upon a worker-pool that allows at most MAX_WORKERS outstanding requests to S3. Not the most efficient solution for this problem. But it works for their use case, for now, and that's what truly matters.
Would something like NodeJs+clusters (or any evented IO framework) be a better fit (considering the clusters are stateless, and don't have to talk to eachother)?
If we're talking concurrency/parallelism, would you prefer JVM+Threads/Erlang+Actors over Go? Thanks.
They DID do a crappy job at concurrency and it's fair enough to admit it. It's not personal sometimes it takes a few iterations to get to the right solution.
I'm working on a contract doing this exact thing at the moment, and unsurprisingly the Ruby code that's in the picture is responsible for just a tiny little fraction of the overall latency.
You should spend most of the time waiting on S3, pretty much no matter what language you choose. If you don't, it's not the fault of whichever language you use.
In fact, S3 performance is so dominant in this type of scenario if you do things properly, that if you want to optimize for speed and can afford to risk data loss (or in our case, the clients poll for confirmation and re-upload data in the rare event of loss), it's generally proving substantially better to write to local disk and do uploading to S3 in the background.
forget about language syntax / compilation complexity.
say i know both very well , now why to go with "GO" ?
thanks
* Fewer lines of code, fewer gotchas and hence easier to reason about and maintain.
* Powerful concurrency primitives (channels, select) built right into the language, rather than a library. A scalable producer-consumer implementation would probably be 100 lines of Go code.
* If your application isn't too latency sensitive (game server, frequency trading etc) then the GC simplifies matters. Its guaranteed to run for a maximum of 10ms out of every 50ms which is good enough for most applications. (but typically runs for around 1ms)
* Some of the tooling around the language is great. There are some great articles (I remember one posted to HN yesterday) about how people wrangled a lot more performance out of their code using the profile tool, for instance.
* Miscellaneous goodies like testing out of the box, an extensive standard library and being able to compile in 1/10th of the time.
An example of a service migrated from C++ to Go - dl.google.com - http://talks.golang.org/2013/oscon-dl.slide#1
These are the benefits I could think of if a programmer knows C++ and Go equally well. However, suppose he has to work with fellow programmers who aren't comfortable with either, Go would be a superior choice. It would take a week to learn most of Go and perhaps a month to grok it. I think C++ takes much, much longer than that to learn properly.
https://news.ycombinator.com/item?id=9844826
In the past, when I try to submit a story, even if it's a couple days old, the submission is ignored, and my vote is added to the original.
Please only do this if you really believe that you should resubmit.
https://www.google.com/search?q=golang%201%20million&rct=j
this HN discussion is too recent to show up in my results, however yesterday's thread about Qihoo and golang does show up as number 4 or 5 when searching for "golang qihoo".
no need to panic.