CockroachDB Stability Post-Mortem: From 1 Node to 100 Nodes
cockroachlabs.com
cockroachlabs.com
FWIW, we had very similar issues at Neo4j a few years back; I kind of wish we'd have blogged about it, but like you insinuate, I guess there's a strong social pressure to keep your dirty laundry in the hamper.
One of the interesting things we learned was that testing "worst case" was often not very useful - the first time we shipped a GA that'd gone through months-long soak testing in the worst possible load we could imagine, it took less than a few weeks for a customer with a seemingly benign load to come back with an unstable cluster.
Today, stability testing is done both on the "max out write throughput, then pull the power cord, then pull the powercord again while it recovers, and do that a million times"-level and "run these four customer use cases, including their devops setup, for two months and fail if uptime is not 100%"-level.
I'm not at Neo anymore, but I do wish someone who is would blog about the stability testing regimen, because at this point, it is brutal.
One of the robustness suites tries as hard as it can to permanently destroy a Neo4j cluster, looking for things like distributed deadlocks, faults in the leader election, data inconsistencies between replicas and so on.
It does that by applying a randomized load of all operations the system supports; reads, writes and schema changes. That's then combined with induced hardware faults and "regular" operational restarts and cluster configuration changes.
The problem was that, early on, the test would actually create unresponsive clusters, but then the "chaos" would continue, stirring up enough dust to get the cluster going again before the "unresponsive" timeouts triggered, causing false green tests.
Hence: Today this suite plays out a "chaos" scenario, but then it heals network partitions, turns unpowered instances back on and so forth, and sits back and waits for the cluster to recover "on its own".
handle a machine power down? Yes. Handle a network cable unplug? Yes. Handle 1 machine power down and 1 network cable unplug at the same time? That may be impossible to handle, or it may be an order of magnitude harder to handle.
Get all the single failure issues nailed down. Then perhaps you can do a little work on multiple unrelated failures, but it gets insanely complex fast.
If I was building this type of system, I would want to build an exhaustive test suite that allowed me to make stability invariant forward progress on the code base without relying on engineers eye catching regressions.
That would mean constructing a system to enumerate all the edge cases at different cluster sizes, automate testing, be able to simulate the network with DI and without docker, and be able configure different network delays and other salient parameters.
The omission of a system like this is surprising because it's really not that hard to build and make it an integral part of core development. I may be wrong, but the last paragraph of the post certainly seems to imply that this isn't done or hasn't been considered.
That is why a simulator is useful. The entire point of building a sim is so you can run the exact code that comprises CockroachDB and control for any situation. For example, an optimization that works great at 10 nodes may be terrible for larger clusters. You need to find that out through testing because most humans can't look at a line of code and infer that.
A side benefit of this approach is that if you change some parameter and suddenly experience a big increase or decrease in db throughput you gain a sense for which parameters are most sensitive to the overall system and will then be able to more easily diagnose problems in production.
(Also: formal proofs of correctness for stuff like leader election are very much concerned with stability / liveness, not just safety; that's a common misconception).
Once they get to the rest of the stability and performance roadmap milestones, it seems like they are positioned to be a great backing store for a lot of interesting use cases.
Nice work!
[1]: https://www.cockroachlabs.com/blog/cockroachdbs-first-join/
I have concerns they will not be able to reach stability in a reasonable time frame and am prepared to pull the plug and go to postgres if needed. The fact that they are trying to maintain compatibility with postgres really makes this gamble far less risky, and was a wise strategic move.
I wish them the best.
I'll keep watching for jepsen tests!! :) https://www.cockroachlabs.com/blog/diy-jepsen-testing-cockro...
[1]http://research.google.com/pubs/pub41344.html
[2]https://www.cockroachlabs.com/docs/frequently-asked-question...
Love @redwood's response though: "This is like a parent saying "We're really looking forward to our child learning to walk by her next birthday" and you responding: "I've got children in their 30s... what kind of parents focus on getting their kids to walk?""
Perhaps somewhat driven by Cockroach Labs home page, which mostly describes the product as they intend it to be, not as it currently is. The github page, blog posts, etc, are very transparent about the current beta status, but the home page is a bit marketing heavy.
Not trying to be overly critical, but rather, thinking it might be the reason for the harsh commentary. (commenters not seeing that the product isn't yet at a 1.0 release)
The alternative is for CockroachDB devs to clearly state "this does not work yet, no one should be trying to use it" before all notable communications.
This is actually a necessary immune response to a "open source ecosystem" that is becoming overwhelming to navigate - people need to know what to not bother paying any attention to yet.
Software is fast paced. You just can't ignore all unstable/beta technologies, you'll be far behind the game when they hit 1.0.
Their reaction to the problems and criticism was good though. They got on the problems. The product is in better confition. A dedicated team is there to keep it that way. Good ending.
But I haven't because of the name.
I want to raise a totally sideways concern that has no relevance to the project, just so it gets out there. The only problem I'm having is that I, like about 10% of the US Population, have Entomophobia.
The name of the project gives me the heebie-jeebies and keeps me from wanting to work with it because I loath cockroaches. I know it's not relevant technologically, but I really can't imagine doing training on cockroach for example.
Oh. My. God. No.
Anyways, wanted to provide the feedback. Sorry if it's totally weird.
I sympathize though. Entomophobia must be a bummer.
There was a time when most software was developed that way.
To build an embarrassingly parallel system; you can't avoid some sort of sharding. I haven't looked too deeply into 'Quiescing Raft', but it still looks like scalability would be limited - It still looks like we have a scenario where each node might communicate with every other node (though in a more economical way).
To remove the limit completely, you'd have to settle with a solution that has fixed-size Raft groups which do not grow as you add more nodes to the cluster.
Since Go code is formatted with tabs, you can mostly get away with setting the indentation to whatever you want. The one practical problem with letting people choose their own values for the width of a tab is that it becomes tricky to enforce a uniform line length, so we've standardized on two-space indents (and 100-char line lengths) across the project.
>Systems like CockroachDB must be tested in a real world setting. However, there is significant overhead to debugging a cluster on AWS. The cycle to develop, deploy, and debug using the cloud is very slow. An incredibly helpful intermediate step is to deploy clusters locally as part of every engineer’s normal development cycle (use multiple processes on different ports). See the allocsim and zerosum tools.
I'd like to know more details about why for this use case, AWS deployment is slower; bandwidth?
udev rule something like: SUBSYSTEM=="net", ACTION=="add", ATTRS{vendor}=="0x8086", ATTRS{device}=="0x10c9", ATTRS{subsystem_vendor}=="0x103c", ATTRS{subsystem_device}=="0x323f", RUN+="/sbin/ethtool -K $name gso off tso off sg off gro off"
Example of an OpenStack vendor guide recommending turning the offloading off: http://clouddoc.stratus.com/1.5.1.0/en-us/Content/Help/P03_S...
Downside is you may have slightly higher CPU usage and/or slightly lower network throughput, I personally have not noticed a significant drop in performance.
Would be cool to see some of the bigger relevant PRs
[1]https://github.com/cockroachdb/cockroach/labels/stability
[2]https://github.com/cockroachdb/cockroach/issues?utf8=%E2%9C%...