325 karma · joined March 17, 2018
I think the worst part is mmap-on-disk looks fine at first, and only comes out as a problem after you scale up a while. False sense of security. :/
Even with an optane pmem mount the perf. is super close.
If you have some kind of pmem device and a dax mount, those can survive reboots on their own.
To find out how long it was down, it notes system time into the state file on shutdown. On start it checks the current system time and adds the delta to the monotonic timer and resumes. Objects exceeding TTL are removed appropriately.
With the restart code, people could run a kernel upgrade and reboot while the daemon is down... so if this ends up causing a huge clock adjustment you're screwed.
The access pattern isn't optimized at all for flash or HDD or etc... however it does work super well if that mount happens to be a DAX mount over persistent memory.
All of the reviews claim it's crud.
At some point I talked to someone who worked at the co who ate the co who would've owned the source. IIRC someone went looking for it but it was gone/lost. There've been some recreations but man it'd be great to get the orignial.
The other k-rad 90's MMO I played, Meridian 59, got its sauce dumped relatively recently too. That was a fun read.
By comparison I've never gotten a complaint about pledge support (or the FreeBSD equivalent), which we also support.
Drooled constantly over the XBAND keyboard that they probably only made ten of. I can still "type" relatively fast with a game controller.
Since they mention memcached: I've been working on a protocol extension to bake this exact thing in more directly. Though in braze's case, it's unclear to me why they didn't use the method of add'ing a secondary key with a low TTL since that doesn't cross systems at least?
With the new protocol to memcached you get "win" tokens, which are very loosely similar to leases. Rather than explicit lease tokens a client is notified of if it "won" or if an object is "stale" etc, and the existing CAS mechanisms are used for replacing objects.
IE: If you fetch an object and miss, it'll auto-create an object with a specified TTL, and return a CAS value (a version number). Winner recaches, other clients are told to retry or wait.
Closer to the braze use case, you can set a "TTL remaining threshold" with a request. If you fetch an object which initially had a 180s TTL, but now has a <90s one, you get a win token. Other clients get the existing value, but only one client is allowed to recache.
There's a bit more to it as the changes are trying to stay flexible for a number of possible scenarios. Cuts out roundtrips and finally gives people more modern cache semantics to work with built in. Hoping to ship this soon, but I need to track down client authors for feedback.
The original question was that it can be really hard to determine how much adding RAM will offload from the disk. IE: with extstore recent objects are served by RAM and never disk, so they don't count against your IO balance. If you're coming from a RAM backed system that has 500k requests per second, it's going to take some creative introspection to figure out how much RAM you'd need to keep the disk accesses below 400k, or whatever your target is.
Along those lines I just suggested a simple experiment. which is to just take an existing fully specced instance, add a disk, and test a few RAM settings to see how it looks. Most actual users of extstore have just done that since they're extending/replacing an existing cluster.
There're actually other ways you can simply read the numbers out of memcached but it just takes more effort or familiarity with its internals.
If you don't have any cache layer at all, that's a standard approach.
2) dunno! For RAM that's super hard. for disk I'm hoping tiering will do something. ie; for on-peak you can add extra NBD space then remove it off-peak.
3) referenced tail latency because most of the testing focuses on displaying latency outliers and how the system generally minimizes them.
heavily batched on a 48 core I can pull 50 million keys/sec over localhost. if you remove syscalls and use it as a library it should double at least.
writes are another story, but they're slower because nobody asks for them to be faster.
1) ID, rack hall, rack number, company name. 2) no id, name, etc. 3) no id, company name, "I think it was left?"
all worked fine.
it ensures that a lot of operations can't touch secondary (like miss, touch, delete, sets of new items, etc), which reduces load on the IO by quite a lot.
edit: Also extstore itself will support NVM + SSD sort-of-layers soon enough. I'll be retesting that on the same optane+ssd machine in a couple weeks.
TL;DR: there's a lot to it and I'll be going into it in future posts. The full extstore docs explain in a lot of detail too.
1) sure, 100b, but that will just make it easier for the CPU version to hit the packet rate limit. I dialed it down to show just how fast the key rate is. Your entire proposal was that CPU bottlenecked the NIC, and it does not. Also, most people have 100b keys, nevermind the values.
2) 1:1 was never realistic. It's not even remotely realistic; as I said earlier 5:1 would be pessimistic. In reality the instances which have get rates in the millions tend to have 100:1 or better ratios due to the nature of the data they're caching.
Yes, the newer LRU algorithm doesn't grab LRU locks on the read path, so it'll scale with the number of CPU cores. As I said in earlier comments, the sets don't currently scale, especially if you're hammering the same LRU (which is again, unrealistic). If you just do a pure set load you'll land somewhere between 900k and 1.5m ops/sec.
3) I did both single-get-pipelined and packet-pipelined benchmarks; also absolutely not. Clients are designed to use the multiget mode when multiple keys are being fetched from the same server. This benefit is lost with the binary protocol (which will be fixed at some point).
4) Try an mget with 16, it won't be too far off, though you might have to add one more mc-crusher thread.
In your last test, you're simply overloading it with sets. If you want to mislead people with a test like this, go ahead; but I'll point it out.
3.5M/s isn't too bad.
Memcached really isn't a great target for your sort of work. I love the idea of FPGA offload, but trying to advertise your thing as superior by making up your own rules is going to get called out.
1) The popularity of redis is absolutely damning in general. if people are okay with the performance of a single CPU database with all-over-the-map latency profiles, the odds of you finding enough customers with extremely high rate memcached pools to sustain a business are essentially zero. You'd be solely tricking people who think they need it.
2) You are not facebook. Nobody is facebook but facebook. 100b is not representative. It's not even representative of facebook's load.
What's worse, even for a more common case, if 99% of requests are 100b, the average size of an item might be 8k. Which doesn't mean that there are a bunch around 8k, but there could be a few thousand items that are 50k-500k+ in size, getting hit 1% of the time, or even 0.1% of the time.
500x the bandwidth of a 100b request for the same processing overhead. It's almost always something they need: a request might fetch a couple hundred items from memcached, with just a couple of them being large.
This ends up making RAM be the greatest expense in the system. If so few users really need this performance, and the newer versions of memcached have a much higher perf ceiling, the extra features it has to drive down RAM usage are more valuable.
The best cost/power savings most users can do is find a way to get more RAM attached to fewer CPU cores: to be frank a r4.4xlarge would suit better with 8 cores. Or find ways push larger cold values into flash, freeing up RAM for those 100b values to be served quickly.
Just signed up for a personal AWS account and manually started an r4.4xlarge for target and c5.4xlarge for source (same CPU's and networking capability?, but it wasn't allowing me to just start two r4.4xlarge...).
got it up to 15M hits/sec for pure mget test.
results: https://gist.github.com/dormando/910134e85279710b970bd2c8af8...
I've plenty of experience with both hardware and virtual machines, I just don't use AWS myself much. I can get something like 800k read ops/sec from a 4core raspberry pi2, and I hope the AWS instance isn't that terrible.
with mc-crusher:
./mc-crusher conf/someconfigfile ipaddress port
https://github.com/memcached/mc-crusher/blob/master/conf/asc... - this is a decent read test with pipelining (give the test a few seconds to get through its sets). The inbound requests are pipelined, but it'll still send each get response in individual packets. This is what I use to test syscall/interrupt overhead.
https://github.com/memcached/mc-crusher/blob/master/conf/mge... this is the same thing, but with mgets. I'd copy the set line from ascii too:
send=ascii_set,recv=blind_read,conns=10,key_prefix=foobar,key_prealloc=0,pipelines=4,stop_after=200000,usleep=1000,value_size=10 send=ascii_mget,recv=blind_read,conns=50,mget_count=50,key_prefix=foobar,key_prealloc=1
can vary the value_size to and mget_count to see how that changes things. You can also pre-warm with the 'bench-warmer' script that comes with it, or remove stop_after and adjust usleep to adjust get/set ratios.
Watch top on the client host, and if mc-crusher is capping out its CPU cores, add more lines to the test but with the (confusing, sorry) threading enabled:
send=ascii_set,recv=blind_read,conns=10,key_prefix=foobar,key_prealloc=0,pipelines=4,stop_after=200000,usleep=1000,value_size=10 send=ascii_mget,recv=blind_read,conns=50,mget_count=50,key_prefix=foobar,key_prealloc=1 send=ascii_mget,recv=blind_read,conns=50,mget_count=50,key_prefix=foobar,key_prealloc=1,thread=1
That puts the first two tests on the "main" thread, then spawns an extra thread for the third test. you can keep copy/pasting that last line until the client or the server are saturated.
edit: sorry, the enc/compression question:
1) compression is typically done in the client to reduce bandwidth overhead. It's not very useful in the server.
2) encryption is becoming more popular, but doesn't currently exist much. The mainline OSS doesn't even have TLS support yet. Almost all use cases are on internal networks. FPGA's could potentially help there... aes-ni on intel cpu's isn't awful though.
Odds are pretty good it's left at the default of 4 worker threads... so on a 16 vcpu instance that's not going to reach great heights. Since it's a 1.4.x version (years old), it's missing some newer features that both help in average latency and memory efficiency. Or rather, a lot of them are there but disabled by default.
Memcached has allowed pipelining since it was created. For the ASCII protocol, packing multiple responses into single packets is done via a straight multiget. You can send multiple requests in a single packet for any protocol and any command.
My stress utility (https://github.com/memcached/mc-crusher) has options for pipelining requests, and using multigets ascii packed get responses. I test to the limit of lock scaling for each individual subsystem.
The 55M test required running mc-crusher via localhost, there's no network that can go that fast. My point is you're limited by the network throughput, not the CPU. In that particular 55M test, all cores were used, but ~7-8 of them were used by mc-crusher... so the real limit for the machine is even higher. It did have a lot of cores. 48ish?
You can still do apples/apples with instance sizes... but given everything I know about this thing, unless those cores are extremely slow, hitting 11m ops/sec shouldn't be an issue. Or at least, with minimal fiddling it should hit 6-8m, which doesn't give you a crazy 9x figure.
You do need to stop doing 1:1 get/set ratio though. Sets don't scale very well since I've generally never had complaints about the speed. I'd say a highly conservative test would be 5:1 get/set. Production workloads are typically even higher than that. (that said I do intend to speed them up more, it's just lowish priority.. the LRU locks are highly granular, so spreading sets across different slab classes can help mutation perf a lot).