Call Me Maybe: MariaDB Galera Cluster
aphyr.com
aphyr.com
See https://github.com/codership/galera/issues/336#issuecomment-... (linked from the article)
Reinforces that the hardest part in engineering is rarely the technical problem. Distributed databases are really f*ing hard, but infinitely harder if the people can't work together or don't open up to faults.
In part it's hard to convince everyone of the importance of spending time on what looks like an unlikely corner case until after the outage - when it also becomes much harder to fix as you usually need to build this into the core design of your system.
It's not hard — it's slow. Sending reads and writes through Paxos or Raft gives you sequential consistency. But not surprisingly, touching a quorum of nodes for every operation is too slow to be practical for many workloads.
And that's usually fine — most data aren't bank account balances.
I suspect that no experienced programmer expects to get data safety for free. We do expect that the documentation doesn't lie to us when describing the level of data safety provided.
What might be considered an appropriate data store for one data set may not be for another. And the type of information a lot of distributed systems are handling are far from bank accounts, for example.
And if you're handling ultra important, sensitive data there are techniques that have been available for many years (two phase commit, for example) that can help.
I love this series but I'm blown away by how many engineers here automatically assume a system is a failure because it doesn't pass a certain type of test from time to time. I do agree with the series that the marketing materials shouldn't claim things, however. Companies should be honest with what their systems weaknesses currently are.
Online transaction processing creates "pending" transactions, and the data is often inconsistent. Your charge may exist in the merchant's database but not post to your online banking for several hours. Or it may be a wildly different amount - i.e. gas stations will place a $100 hold on your debit card and it will stay that way for days until it's settled for the actual (lesser) amount. If you were accidentally double-charged, than rather than processing a separate refund, the merchant may simply not settle the duplicate charge, and it will drop off your pending transactions... eventually. The lag time may be several days or weeks. If it's a debit card, you can't spend the money during that time and you may be temporarily broke because of it.
If you make an ACH transfer, money will disappear from your account one night, spend a business day in the aether, and then post to the recipient's account on the third day. The system is in inconsistent state (i.e. money is missing) for at least a day, possibly a whole weekend.
The actual transaction is settled and goes into "posted" state with a lag of 2.5 * 10^8 ms - i.e. 3 business days. That's if you're lucky. Banks do need strong consistency, but not in anything approaching realtime.
Even ZooKeeper could probably handle the U.S.'s financial transactions faster than current infrastructure.
Probably not even going via ACH, rather through treasury department nettings, mitigated against each other, just because that's accepted.
RTGS really is Real Time Gross Settlement. If you're led to believe elsewise, you're being conned.
There will be bugs but in the end if your goal is to write a consistent store, you can do it. Bugs in wanna-be CP systems and/or failed attempts at implementing CP because of home-made protocols with flaws are part of the issue and part of what Aphyr's research uncovers, but IMHO the central point is a different one: CP systems have performance limits so many real world systems don't even try to be CP, which is fine, it's up to the designer to pick the DB design and tradeoffs, but you need to document it properly. What happens during partitions? What consistency model the system employs if it does not feature strong consistency? Is the system Available during partitions? And so forth. As long as the documentation is honest, it's up to the user to understand, for its use case, if the system is a good fit or not.
What we are seeing often is an attempt to informally document what should be well specified. Sometimes there are also incorrect things stated in bold letters in the documentation, like "the system is consistent but sacrifices partition tolerance" which does not mean anything in the context of partitions actually happening in the real world.
Aphyr's effort and our collective effort as people working, implementing and using databases, in my opinion, is not to end with databases designed all the same way or all providing strong guarantees, but with clear understandable and honest documentation describing the system behavior.
They try on paper, so to speak. They might have a CP component, but it is not in the datapath -- it is there for maybe cluster configuration, or say to elect a leader. And then they can claim they have CP because everything goes through the leader, and the leader was elected by a CP component. However without realizing they just pushed all the error cases into corners -- when the membership changes (nodes die, die in unpredictable ways, become partitioned, new nodes added...). They never test that extensively, which is a hard thing to do, because there are so many cases and combinations.
The bottom line is even if they have a Paxos or Raft component in their system doesn't mean their data is stored consistently.
I am a MariaDB contributor and operate Percona Cluster in production, so I can talk a bit about Galera.
It's recommended that writes go to one master, rather than be distributed across the nodes. That will help with isolation issues.
Also, some commenters have complained about year-old releases. PXC has improved significantly in the past year regarding manageability, so you may want to try again. For example, the startup script has a bootstrap option now.
For most people, vanilla async MySQL replication works best, esp. 5.6. But Galera gives you another option when you need something else.
Having said that, it takes 5-10 years for a database or filesystem to mature, so anybody using Galera now is an early adopter.
Essentially, what a pain in the backside to use.
This did not fill me with confidence and thankfully did not go into production and later on went with Postgres/CitusDB.
The difference is day and night!
I've been strongly considering MariaDB on Galera for encryption-at-rest. Is there something about MariaDB that was not working well with Galera?
And why noone tried to repeat their strategies for building a robust db system : start by building an extremely robust failure simulation and testing facility. Then build your product.
Actually, i think what those guys at foundationdb did was so exceptional, that by buying the company and killing the product, Apple harmed the software industry for the next 10 years. The fact the foundationdb is mentionned in OP as the only distributed db system one could recommend makes me more confident making that statement.
"NoSQL" is not a magic bullet: concurrency is still hard when you skip the SQL.
I meant that with this being the first distributed SQL store, he was saying there was nothing he could recommend that actually offered similar guarantees to what this was claiming. That is, ACID style (well, technically --ID I believe) DB transactions in a distributed environment. FoundationDB did (though it was NoSQL), hence his mentioning it. But that's different than there is nothing he can recommend, at large; he can't recommend similar solutions, but for a given problem he could likely recommend a compelling, different, solution.
(What follows is my opinion / guesswork as a VoltDB employee)
1. FoundationDB wasn't killed by Apple; it was rescued by Apple. The product couldn't compete on just being a KV store and wasn't doing well in the market. Apple saw a very bright and now experienced team and scooped them up for a song.
2. Before this happened, FoundationDB realized they needed a way to query their system to compete, so they bought Akiban (a failing SQL db company) to add SQL to their system. But they assumed they could do this without deep integration, which was wrong. They added a SQL "layer" on top of the KV store and it was way to slow to be practical. The benchmarks they published were embarrassing.
I wrote a blog post about this: http://voltdb.com/blog/foundationdbs-lesson-fast-key-value-s...
SUMMARY FoundationDB: Great Testing, Great Engineering, Not particularly good product...
Actually, the thing that i find most impressive in their tech stack is the approach they took for building it starting with the simulator + c++ extension. Those are the technology that i think would benefit all the community, if they were ever open sourced.
As a voltDB engineer, how do you ensure your implementation doesn't compromise the theorical correctness of your system ?
http://www.thestrangeloop.com/2015/all-in-with-determinism-f...
VoltDB is actually a bit simpler with what it promises, full serializable ACID for all transactions. This is much easier to understand and verify than lesser isolation.
We think what we've done is pretty clever too. We've built a determinism checker into our replication engine, so that we can verify that each replica has the same state at each logical point in time, and each operation makes identical changes to that state.
Then we built test patterns that are designed to be as co-dependent as possible and run them against a replicated VoltDB cluster. That VoltDB cluster goes though one or more kinds of failure, including multiple simultaneous failure, and then a checker ensures no data is ever lost, corrupted or run in the wrong order.
It's different from the FDB thing. The simulation they did is certainly easier to run on a pure KV store, but keep in mind we also have to test SQL that queries millions of rows, along with upserts, materialized aggregations, etc...
We're working on some blog content on this in addition to my talk. Stay tuned.
But, since i'm absolutely not working on the field, i'm really looking forward to see what professional people like you are finding to tackle those issues.
If you're talking about incentives, don't forget that Kyle's research is some of the most commercially useful research out there at the moment, and Kyle is also skilled at implementing his ideas (a rare combination among researchers). I'm actually quite pleasantly surprised that he's putting so much effort into these blog posts rather than writing a string of journal papers.
Sure, I agree (well, I'd probably replace "much" with a less strong word, and skip the "damning", but that's just semantics). Plenty of studies haven't been demonstrated to be wrong, and plenty of others were wrong despite being correctly designed experiments (something both peer reviewers and the scientists performing the experiments could not have caught). Moreover, many flawed studies haven't been accepted by peer reviewers, which is in fact the purpose of peer review. The biggest problem with peer review is probably publication bias against negative results, without which I suspect most demonstrated scientific fraud wouldn't exist, but that doesn't mean the whole process needs to be thrown out.
> Further, I'd say that with the tools being provided for free and the blog posts Kyle's findings are most likely peer reviewed... And copany reviewed... And product team reviewed..
Yes, Kyle does (hell, he's been cited in academic papers). The average person publishing blog posts does not, though.
Sharding is good, for sure. It's not a replacement for this kind of technology.
Aphyr is doing great work keeping vendors honest, and making sure they live up to their marketing. I don't think any of this says that clustered databases are inherently shit.
"You might adopt a different database–though since Galera is the first distributed SQL system I’ve analyzed, and FoundationDB disapparated, I’m not sure what to recommend yet."
They ran some tests themselves using the Jepsen framework, but given that a large part of the testing is setting up an environment that exposes the system's weaknesses, that doesn't give me the same level of confidence.
https://aphyr.com/posts/282-call-me-maybe-postgres Also seemed solid
It is irrelevant in this discussion.
But regardless you completely missed the point. PostgreSQL was tested as a single instance. Everything else clustered. It's apples and oranges. The clustering part is the hard part here.
My point was simply to remind people that over complicating your situation and going for a "big data" store is usually a mistake when there's something that will keep your data safe when it isn't as big as you assume it's going to be.
Who'll ever know if a search result isn't perfectly up to date or perfectly accurate?
Who'll ever know if you missed a Facebook feed entry because it "wasn't relevant" or simply wasn't seen due to DB vagaries?
And who'll ever know about a few tweets going astray here or there?
In all cases, they're all likely "eventually consistent" (or close to it), but it's no accident that it doesn't ultimately matter in those massive scale examples.
And maybe that's the secret to massive scale--it can't ultimately matter.
We basically queued all retrieved items for processing with no attempts at avoiding data loss whatsoever - including using in memory queues for lots of things.
Our reasoning was that if a machine crashed, worst case was that a few listings would take up to 24 hours to update, but generally much less (we adapted crawling rate for our sources based on change frequency; so large sources of listings would get re-indexed far more frequently, so if a feed didn't update or 24 hours it'd be because it wasn't a source of much data), and we could force refreshes of the data.
Some people were horrified at the approach because the idea of ensuring consistency and not losing data is so ingrained. But the reality is that you need to measure the cost of consistency up against the value it provides. And often it's not very valuable, especially when there is an authoritative source of the data to recover from and when the data will be outdated quickly anyway.
A lot of the time any notion of consistency is an illusion anyway - by the time the page has returned, the results are outdated - and what matters is maintaining the illusion (e.g. ensure that if a user makes an update it's reflected in the page that's returned).
The key is you need to know the tradeoffs and apply them consciously rather than get caught out by tradeoffs components you rely on makes without telling you..
Game servers use the same technique.
If I had gone to his lengths to critically evaluate the safety of a database system, and then it comes out that the marketing materials or the words of the developers were.. significantly misleading at best, my first response is likely to be a profanity laden rant, not a cool recounting of how and why they're wrong.
Its not so much about wanting to bolt onto amazon's eco for stuff -- but just allowing an org to focus.
AWS allows for an ops team to completely have no concern for hardware. Lovely.
I also don't want an ops team doing DB cluster mgmt.
I want to deploy and delve into my data.
I want my data/eng/dev/ops teams to be pummeling the shit out of my data without concern for the instances/hw/cluster.
aurora adds to the cloud fabric.
One thing amazon is moving to is to define compute capacity vs defining instances at all.
Their vision of abstraction is amazing.
I am sure this has been looked at by people -- anyone know of a report on it?
The company I manage physical servers for would go bankrupt if we were to pay AWS level rates. Last time we priced it out we were looking at ~3x hosting cost - including fully loaded cost for staff time spent on maintenance etc.
I split my time between managing a few racks worth of servers at one company, and an AWS setup at another client. The amount of ops work per server/instance is higher for the AWS client than it is for the company I manage physical racks for. That's despite the fact I do physical cabling and racking of equipment as necessary.
The issue is that there are so many aspects of the AWS infrastructure that takes extra effort because we don't have full control. E.g. we can stuff whatever disk subsystem we want in the servers and not have to work our way around the lack of any truly fast disk subsystems in EC2. So for every day I don't physically move serves, I spend 3-5 working around AWS limitations.
I agree with you with respect to what you want, but we're taking baby-steps today - AWS is way too expensive for most people to move to it, even before you factor in the ops complexities.