Russ’ 10 Ingredient Recipe for Making 1 Million TPS on $5K Hardware
highscalability.com
highscalability.com
You could get rid of this (and in doing so, double your TPS) by switching to a memory-polling-based network driver like PF_RING [1] (and obviously, keeping the kernel on its own core like you are doing).
> Lookup data in memory (this is fast enough to happen in-thread)
> I knew I had it right when I watched the output of the “top” command,
Does all your data live in cache? If not, have you tried using perf [2] to measure load stalls? They are typically the bottleneck once you get past context switches. Hyperthreading should help at least somewhat here (do you have it enabled?).
I am surprised (from a quick Google) there is no open-source user-space PF_RING-aware TCP stack. Am I missing something?
It reached an awesome Kps - keywords per sentence count as well.
1 million network-to-memory writes, well, that is quite possible, but, please, do not call this a transaction in the way it meant in TPS.)
What was meant in the old days by transaction, was an atomic operation which completes after storing the data in a persistent (usually direct-access, which means no buffering by an OS kernel) storage, so it could be read without any corruption if a power loss will occur the very next second.
In any case what they seem to be measuring is a read-only, in memory workload, so this is not that impressive. IIRC, InnoDB has no trouble pulling off something like this.
Yes, there are lots of tricks, like placing that append-only physical transaction log on a different controller with a distinct storage device, etc. Data partitioning is the another big idea. Having indexes in memory to avoid unnecessary reads, using collected statistics in a query optimizer, etc. But nothing could beat the partitioning based on actual workloads and separation of tablespaces on distinct hardware, including decoupling indexes from the tables - this is what DBAs were for.
I used to be Informix DBA in old good days, so I can't help but smile when I look at MySQL (well, they added lots of partitioning options in recent InnoDB - the things Informix could do out of box 12 years ago) leave alone modern NoFsync "databases".)
Btw, not all NoSQL guys are insane.) Riak with LevelDB storage backend is very sane approach, which cares about and counts writes.
Partitioning is often a bad idea because it messes with your queries. I don't know what's new about partitioning in InnoDB but I think it's generally a symptom of the over-use of B-trees, which don't try to do anything smart about random writes. The change buffer is a decent idea but it's just a stopgap, when you have enough data it doesn't make a dent any more. A better idea is to use a data structure that can handle lots of writes.
I have just started learning about Riak, and from what I understand, they need to do a query (so, a disk seek) on every insert (to calculate something with vector clocks), so they aren't actually using the write optimization that LevelDB's LSM-trees can provide. I don't actually think it should provide that fantastic performance, but I should admit I haven't run it yet. Maybe they're more interested in the compression LevelDB gives them.
Shameless plug time! http://www.tokutek.com/2011/09/write-optimization-myths-comp...
It should get very interesting in the next couple of years.. of course MOST environments don't need the kind of performance or scale that these systems are really offering.
IIRC StackOverflow ran for a very long time on a single server, under some pretty serious demand. In some cases SQL with a caching system for mostly-read data can be better... other scenarios tend to fit a document (non-relational) data store better.. just depends.
You're right, with a reliably performant engine you can get a lot more out of a single machine than a lot of people these days seem to think. That's part of our vision for TokuMX, to bring back a little bit of "scale up" potential to the NoSQL space.
They also discuss how these are 1M read requests.
If one would like to delve a bit deeper into this one can find some info regarding the techniques introduced for linux here: http://lxr.linux.no/linux/Documentation/networking/scaling.t...
Does anyone have some actual configuration examples to provide for things like setting IRQ affinity?
(I single EVE out because other MMOs generally shard by user-cohort, so having that number of people on one shard is impossible. EVE, meanwhile, shards by location within the virtual world (each star system is a shard), so the entire player-base can "gather" on a single shard for a confrontation.)
So what, the entire game userbase needs to be awake and in the same spot? :0)
Specifically during DDoS attacks, an IPS must usefully distinguish bad from good traffic at packet rates saturating a link. This inevitably involves maintaining lots of per-client and per-IP state. Caching is of little help precisely due to the distributed nature of the attack, and you can't shed load since that only helps the attacker.
In this situation, memory stalls become your biggest bottleneck – each one can eat on the order of 10% of your processing budget in a run-to-completion (RTC) design. The only solution (beyond tricksier data layouts) is memory latency hiding via micro- or hyperthreading (kernel context switches are just too slow). Rearchitecting a RTC design into a micro-threaded model is a lot of work, and bug-prone. Hyperthreading gives you latency hiding "for free", if the silicon supports it.
Generally you want between 2-4 micro- or hyperthreads. A second micro/hyperthread will generally just help keep the pipeline busy outside memory stalls; hence you can eke out extra performance with a third or a fourth. Intel chips only support two hyperthreads (when they do, and the OS supports it). Some more specialized processors support more.
For example, I've had some issues with CPU (on core 0) saturation for soft interrupts on EC2, related to processing network packets.
Some hybrid stores like MongoDB, RethinkDB and others offer more characteristics similar a traditional SQL RDBMS, while offering horizontal scaling. It's when you have several join operations that performance really takes a hit under significant load. You can't really scale a relational database in the same way.
That said, as much as I like non-relational databases, they aren't the best fit for every use case. Beyond this, most situations don't need that type of performance scaling.
If only they would do the reader the simple, and most gracious service of defining precisely what they mean by this obscure "TPS" acronym.
...and before you downvote this comment (because I can smell your itchy little fingers all the way from the otherside of the internet), yes, I can assure you that I did actually Google for the answer. And yes, I did discern what is meant by TPS.
But that isn't the point. The point isn't that I'm a lazy slacker, and/or an ignorant yokle because I didn't already know the meaning of the abreviation innately, and feel inconvenienced by having to open another browser window, and search for some clue.
The point is that the author is assuming everyone will immediately know and understand that acronym, but meanwhile, when I conduct my search, I am forced to assume that my chosen definition is correct, wihout actually knowing for sure.
And for that reason, I'm going to leave out the meaning I've chosen as the author's intended definition for TPS. I have no way of knowing whether my assumption was accurate. So the mystery persists. What does TPS actually mean? Go search for it, you lazy, ignorant slacker.
If you don't know what the acronym means, then the article is not meant for you. The author did not write it nor intend it as a general introduction for newcomers to learn the basics.
It's about some some specific techniques in a specific field.
With your logic, why stop at explaining TPS? He would also have to spell out IRQ, explain what IRQ interupts are, what is a tasket, what does it mean to 'pin' a process, what NoSQL means, what's this Redis thing he mentions etc etc.
(That the author of TFA abuses the acronym is beside the point).
A certain level of reader sophistication is assumed, and not everything is explained with a high level of hand-holding. It's a technical tutorial, not a beginner's tutorial.
TPS = transactions per second.