84% of a single-threaded 1KB write in Redis is spent in the kernel
blog.nullspace.io
blog.nullspace.io
This is relatively common for closed source software but you almost never see these types of userspace I/O designs in open source, which means OSS designs are often leaving integer factors worth of efficiency and performance on the table. For some use cases, companies are very successful selling into this efficiency and performance gap with closed source.
Part of the lack of open source is that these designs are not portable, due to OS dependencies and sometimes hardware dependencies if a design is hardcore. I think this is a partial copout; a Linux-only database engine would address the vast majority of real-world deployments. A bigger reason is that the design and implementation of these kinds of userspace kernels is a very high skill and low level art that, frankly, is way outside the expertise of most open source software contributors. For databases in particular, more often than not even the basic design elements of the software are naively done (e.g. MongoDB) and that is lower hanging fruit.
Basically to make Redis much faster we need to work to three different related things:
1) Less kernel friction.
2) Threaded I/O, this is the part worth threading, with a global lock to execute queries so you don't get crazy with concurrency and complex data structures. Memcached did it right.
3) Pipelining: better client support for pipelining so that it's easy to tell the client in what order I need my replies, and unrelated replies can be glued together easily.
http://symas.com/mdb/memcache/ http://symas.com/mdb/inmem/ http://symas.com/mdb/ondisk/
This is one of the big reasons I like the paper I mention in this post.
So I'd be willing to say that the problem here isn't that the kernel stack is slow per-se, but that workload is too small as to make the overhead look ridiculous, when it'd be very much acceptable if your server did more actual work.
Performance is usually a secondary concern in these use-cases.
And all of the Redis set/list operations are super valuable. You just need to be careful once you start relying on Redis at scale for things you take for granted when you start playing around with it... For example: zunionstores on lots of large sorted sets. "O(N)+O(M log(M)) with N being the sum of the sizes of the input sorted sets, and M being the number of elements in the resulting sorted set." When you start off using it, it's awesome and super fast, but before you know it, the blocking, single-threaded architecture will crash and burn if your data scales up. Luckily we have clustering now :)
Regardless, your goal should be reducing shared state as much as possible, since synchronizing it among distributed systems is a Hard Problem(TM).
Also, saying that using Redis is slower than using a local hash table is a truism. There are myriad reasons why using a local, in-memory data structure is not viable: scalability and persistence, for example. It's like saying "I don't need a database, I can store everything in a local variable."
You may not have these problems, in which case in addition to "the concept of redis" baffling you, this will seem absurdly performance-sensitive to you. From a desktop programmer or all but the most complicated websites, that is also a sensible perspective. But the niche in which this discussion makes perfect sense is itself pretty large.
alongside very aged processors. The hardware was almost 10 years old in both cases. Even the current gen aren't particularly powerful, coming in at 8GB ram with a 1.75GHz processor, and a GPU comparable to a 3-4 year old PC for the xbox one, and 1.6GHz processors, 8GB ram and a slightly beefier GPU in the case of the PS4.
(By "done right!" I mean either batching requests to/from the kernel, or using a zero-copy userspace solution like DPDK. Clearly in this article network I/O was not done right. Round trips through the kernel ALWAYS will kill performance.)
Redis has a great position as a persistent, shareable, data structure server, but replacing in memory hashtables where they work doesn't seem like one of those cases.
Of course, you can pipeline memory accesses to some degree, but not as easily as you can aggregate network requests into fewer packets.
I'm certainly not saying the network overhead is free -- it's not! I'm just saying it needn't "eclipse" the hash table lookup itself (as the GP suggested). They're on the same order of magnitude.
Also, on the target server, you still have at least one cache miss to retrieve the item from the hashtable (and potentially more for large hashtables).
It would be nice if there was a commonly supported API for doing this kind of stuff that doesn't require exotic hardware and came with a usable abstraction something like TCP. I know off-the-shelf Intel NICs have "Direct NIC Access" which can easily do wireline work from userspace, but afaik that's only for Intel NICs, and the API isn't as smooth as most socket users are used to. We need something like SuperSockets, but written to target NIC hardware directly, I guess.
It doesn't make sense for each app server to have an 8gb hash table in their memory nor can they easily be synchronized.
This does tie you in to a specific network vendor but I can see the argument for moving more commodity networking hardware in this direction too.
Improving Linux networking performance
https://lwn.net/Articles/629155/
HN discussion: https://news.ycombinator.com/item?id=8931431
There's been periodic interest in 'virtual hardware' where the hardware presents multiple interfaces to different users. This way the driver can run in user mode, since there's no need to control/share the hardware registers in the kernel.
Are you sure?
Of the total 3.36 μs (see Table 1) spent processing each packet in Linux, nearly 70% is spent in the network stack
If it is in the IP/TCP layers then moving that to user-space does not, by itself, necessarily reduce latency, it merely shifts it elsewhere. If the latency is due to kernel user land memory copies then that is a different matter.
Consider, for example that the dominant costs are things like demultiplexing and security checks. If you choose to implement multiplexing with virtual network cards then you get true 0-copy multiplexing, which is much faster than the software equivalent. And many of the security checks can be eliminated by using some combination of packet filters and logical disks. (The security BTW seems to be one big difference from RDMA, which might be an alternative, but I'm not really an expert.)
Some things can't be sourced to the hardware, like naming and access control. But that's fine.
(NB, I'm not arguing for this paper's position necessarily, I just thought it was interesting, and the motivation was good enough to start me thinking about how I might get around the kernel.)
Pushing more of the stack into hardware is probably a good idea for single-tenant datacenters that can deploy a lot of e.g. Redis appliances, but those of us just renting capacity in the cloud are going to suffer from Amdahl's Law if you can only accelerate the part of the system adjacent to real hardware NICs.
The kernel socket API is designed so that programs have to do as little thinking as possible to get their own personal slice of the shared and noisy network. It provides an easy abstraction, and that requires the kernel do a lot of messy stuff for you:
- When you're using TCP sockets, the kernel makes copies of everything your application writes and holds it in a buffer until its receipt is acknowledged, just in case it needs to resend it when the other side doesn't acknowledge it. If the socket's buffer fills up, your application blocks on I/O until some space is freed.
- It holds ports open in a lingering state long after they're closed just in case it needs to re-transmit the last bytes. This can be disabled, but it's on by default.
- It takes care of all the congestion control for you, but it's tuned for the general case, and as a result there are a lot of edge cases which perform very badly for the problem they're trying to solve. Redis is probably one such edge case.
Of course, all of this is fine and desirable for general applications, but it ends up being problematic if you're trying to solve a problem where performance is the chief concern.
It's tempting to say the problem is that kernel has to do way too much to provide that easy abstraction, but really the problem is that the kernel provides no way around it. You pretty much have the option of using their cushy stream abstraction at the cost of performance, or you use a userspace TCP stack on raw sockets, which requires running as root and disabling TCP in the kernel (otherwise the kernel stomps all over your TCP negotiations[1]).
There are some other transport layer protocols (SCTP, DCCP, etc.), as well as application layer protocols built on UDP, that remove some of the abstractions TCP provides and as a result require less in-kernel bureaucracy, but those solutions don't seem to be very popular or well-supported.
It would be nice if the kernel would provide some lower level system calls that could be selectively used to move parts of TCP into the application (e.g., retaining copies of data in case of re-transmission). Alas, I don't think there's much push for that, because a) it's hard, and b) the current situation is fine for 99% of network applications.
[1] http://jvns.ca/blog/2014/08/12/what-happens-if-you-write-a-t...
"How did you measure the time spent in each section (HW, kernel, app)? how did you get such granularity?"
I just whipped this up in 5 minutes and didn't do much of the tuning there (e.g. no isolcpus or interrupt changes), but here's a single-client 1024-byte SET redis-benchmark running against localhost with and without TCP Loopback Acceleration... redis 2.8.4 on a dual E5-2630 @ 2.30GHz, card is SFN5122F but this is all loopback. I'm not claiming anything and just doing it because somebody pondered...
* plain jane
/usr/bin/redis-server
redis-benchmark -t set -q -n 1000000 -d 1024 -c 1
SET: 21258.96 requests per second
* unaccelerated server, unaccelerated client
numactl --physcpubind 1,3,5 --preferred 1 /usr/bin/redis-server
numactl --physcpubind=7,9 --preferred 1 redis-benchmark -t set -q -n 100000 -d 1024 -c 1
SET: 14293.88 requests per second
* TCP loopback accelerated server, unaccelerated client
EF_NAME=hn EF_TCP_SERVER_LOOPBACK=2 EF_TCP_CLIENT_LOOPBACK=2 onload -p latency numactl --physcpubind 1,3,5 --preferred 1 /usr/bin/redis-server
numactl --physcpubind=7,9 --preferred 1 redis-benchmark -t set -q -n 100000 -d 1024 -c 1
SET: 25967.28 reque sts per second
* TCP loopback accelerated server, accelerated client
EF_NAME=hn EF_TCP_SERVER_LOOPBACK=2 EF_TCP_CLIENT_LOOPBACK=2 onload -p latency numactl --physcpubind 1,3,5 --preferred 1 /usr/bin/redis-server
EF_NAME=hn onload -p latency numactl --physcpubind=7 --preferred 1 redis-benchmark -t set -q -n 1000000 - d 1024 -c 1
oo:redis-benchmark[13454]: Sharing OpenOnload 201405-u1 Copyright 2006-2012 Solarflare Communications, 2002-2005 Level 5 Networks [4,hn]
SET: 96098.41 requests per second
Edit: formatting fixes