3 ms·
I like Slide 251: 'solve it at the application layer' :) I'm part of a team responsible for cloud-based video editing software. We use multi-master replication
by jermy 15y ago
I like Slide 251: 'solve it at the application layer' :)
I'm part of a team responsible for cloud-based video editing software. We use multi-master replication (perhaps also known as optimistic replication) with our own tools throughout, but it does require careful design - keep as much data as immutable as possible, every piece of data that might be updated by different machines at the same time should have its own row, GUIDs on each row.
Each machine can generate its own local IDs, which look a lot like a timestamp with some unique stuff on the end. Each row gets a GUID and a 'version' ID column, and we only update relayed database updates if the incoming version is newer. This is largely last-timestamp-wins for the case of conflicts (rare because of design decisions), but there is some Lamport timestamp behaviour in there too for updating a existing row.
The main downside is still that all machines need to handle every write, but with batching up incoming processing into larger transactions, we've had no problems with quite a number of database updates on a dozen commodity machines. Obviously filtering into different shards would be an easy solution.
I'm looking forward to seeing what other people are doing with multi-master replication.
- spydum 15y agoPretty typical of any multi-master asynchronous replication: you must make damn sure your application isn't going to generate conflict. If conflict does happen, you have to apply some external logic to resolve that issue. Some things like GoldenGate (oracles purchase of binlog snarfing multi-master replication) provide you some tools for that. Basically dumping conflicted records into a table for manual or automated clean up later. At the end of the day, it's always going to require well thought out schemas and data layout if you want to take writes for both masters. Multi-master is less difficult if you only write to a single node at a time (hot-standby style).
- alfiejohn_ 15y agoHaving a single master is sometimes not viable in some situations. But yes, eventual consistency can be both a blessing and a curse. But if you can design your system to take this into account, it's pretty easy to scale out and everything Just Works(tm).
- mattw 15y agoI don't know much about replication, but I'd like to learn. Can you recommend any books/courses/tutorials that address the sorts of design decisions you're talking about?
- alfiejohn_ 15y agoThe two must have books for MySQL admin and discuss MySQL's built in replication, what problems you'll encounter, how to tune etc. - High Performance MySQL (Schwartz et al) - MySQL High Availability (Bell et al) I think using MySQL::Replication can be applied on top of what's in these books, but not using their particular setups.
- alfiejohn_ 15y agoToby will be uploading the talk to Vimeo when he weekly upload limit gets reset. We talked a bit about how race conditions are a part of life with multi-master replication, and the ways to best avoid it. When you know your data, you know what you can get away with.
- IgorPartola 15y agoAt TransLoc we use a multi-master setup. At its heart it's two nodes replicating from each other. Our data is divided into several databases. Each database "belongs" to only one master at a time. For example if we have nodes A and B, and databases w, x, y and z, we would have A be responsible for writes to w and x, while B would be responsible for y and z. If B fails, we have a monitoring system in place to tell A that it is responsible for w, x, y and z at once. The monitor sets read-write permissions for databases y and z on A (through user permissions), and then notifies all of our application servers that things have shifted. The applications for their part include a piece of common code that monitors for changes in the cluster and allows the application code to cope with these changes. For web requests, if failover happened in the middle of a transaction, the request fails. For long running processes, the process will have to go to the top of the event loop and request a new database connection, etc. So far it's worked fairly well. We are able to achieve high availability with it, since our master nodes are in two different data centers. There are definitely issues with this approach in general, but it works for our work load.
- alfiejohn_ 15y agoYou're kind of in an active-passive multi-master setup which is find for now but as replication queries increase between the two nodes, you might start to see load issues. You're going to have to drop in a third machine and a decision is going to have to be made on how to replicate your data.
- IgorPartola 15y agoYes. The other big problem with this setup is that a single node must be able to handle all queries if the other node fails. Fortunately, we are nowhere near saturating the capacity of our nodes under normal operation (we are only doing about 1,200 queries/second). We are also able to offload read-only queries to slaves that are replicating off of the masters, which should help with the read capacity, which is our biggest demand.