Pure Operation-Based Replicated Data Types
arxiv.org
arxiv.org
In particular, I've wondered how a simpler solution to the problems posed in this paper stacks up. At what point does it fail? e.g., Consider:
* An op based model.
* A client-server architecture (not p2p).
* Operations are assigned a total ordering by a central server.
* Given any two dependent operations A and B, A is always
transmitted before B.
This still gives you a lot of the stated advantages of an eventually consistent system, where each client communicating with the server will eventually converge even if they all temporarily diverge. The central server and total ordering is a key ingredient, because it's what lets you guarantee ordering between any two causal operations by having the server "choose" an ordering.I'm naturally interested in trade offs. For what use cases does the central server model break down? Is it still useful for other things at lesser scales?
[1] - https://medium.com/@raphlinus/towards-a-unified-theory-of-op...
I work on a realtime database and we heavily rely on CRDTs for communication. We can't have a central server because our goals are fault tolerance and high throughout.
If you're more curious about the subject look into Riak. Afaik that is the biggest project built with CRDTs
I guess I'm trying to figure out why this is a problem. If operation B shows up hours later, then the client simply doesn't know about B until then. But since my stated design requires that dependent operations are delivered in order, we have that any other operations are necessarily concurrent with B, so it shouldn't matter when B is "evaluated."
> I work on a realtime database and we heavily rely on CRDTs for communication. We can't have a central server because our goals are fault tolerance and high throughout.
> If you're more curious about the subject look into Riak. Afaik that is the biggest project built with CRDTs
Yeah, I'm definitely aware of Riak. I guess I'm more interested in the systems layered on top of these components. i.e., What problems are they trying to solve?
For instance, what if you could relax your requirements? Maybe: throughput/latency only need to be good enough for web response time to end users.
When I worked with CRDTs---even inventing a new one based on R*-trees---they were most useful for objects/structures that needed constant mutating by lots and lots of participants concurrently where it made sense to trade constant coordination/protection overhead for lazy-evaluation/merging. Like in a settlement ledger, or massive leaderboard, or indexing structure, etc.
I'm not sure I see the additional complexity. With a central server, you get to enforce a total ordering automatically across all operations.
In any case, what I'm looking for is actual data. I implicitly understand that a central server won't scale as well, but that's not particularly interesting to me. For example, what if my requirements on throughput are very low? I want to know the breaking point. So I guess I'm looking for an experience report from someone who has actually gone through this process.
(And yes, I've been through this)
A client could be disconnected from the network for hours and re-synchronize at some later point in time and it all would work just fine.
The only time the client needs to connect to the server are when the client wants to publish updates, or when the client wants to receive the updates. If you only need to do that once every 10 seconds (or even once every second), does that make a client/server model workable?
To make it work that way a client has to ask the server for a starting point before it can generate operations and increment it with each new operation. But then it doesn't sound like the server does the ordering and starts to look a lot like lamport timestamps.
For the client, the only real choice they have (until they hear otherwise) is to assume that all local operations come after all server operations. At some point, the server will send the client updates with new operations from other clients, or perhaps even confirmation that some of their previously local operations have now been recorded by the server and assigned a total ordering. The client is then responsible for updating its local state: moving the local operations into the server operation list and updating the server operation list to contain any additional operations reported by the server. The client must then re-materialize any data structures dependent on these operations, but depending on the nature of your operations, this might be possible to do efficiently.
Given enough time where every client is idle, its set of local operations will become empty as they all move into the server's set of operations. All clients then have the same set of operations in the same order, which gives you consistency.
There is still the issue of a node becoming hot because there is an unusually active document on that node, which would usually not happen in a completely decentralized model.
Operation based CRDTs do seem more limiting compared to State Based though, so you're at the mercy of your data structure being able to be represented as a CRDT.
> In addition to being resilient to network split brain, this allows you to make performance gains at the sacrifice of consistency, where you could buffer your updates to other nodes.
My scheme allows this too, no?
I cannot imagine designing a system for production use where anyone is OK with a single point of failure being a single instance of a machine. I mean if we're playing hypothetical games you don't need a distributed system either, just get one giant instance and put everything on it.
> I mean if we're playing hypothetical games you don't need a distributed system either, just get one giant instance and put everything on it.
I guess the requirements need to be more explicitly stated. You need to at least be able to support multiple simultaneous clients interacting with the system, where any client could disconnect for some period of time.
Distributed version control relaxes those constraints with better merging to ensure commits are commutative. Now take that property and extend it to general programming, and you have similar advantages with eventual consistency.
It's harder to see with typical concurrent abstractions, but Microsoft's concurrent revisions explicitly models their concurrency abstractions around distributed version control. Their papers are worth a read if you want ideas, data and inspiration that concurrency and distributed programming can actually be easy.
In any case, thanks for the pointers about the Microsoft research. I don't think I've seen that. I'll check it out.
Each client can operate on its own what though? If they have some kind of private copy of the data, then how do you resolve conflicts? CRDTs are about doing this automatically by restricting the types of changes clients can make. Concurrent revisions permit arbitrary changes so long as you provide a merge function to make them commute.
I think you're just hand-waving a lot of the complexity away by using "synchronize" without explaining what is synchronized and how. That's where all the complexity in distributed computation stems from.
Imagine the fundamental representation of your data structure is a sequence of operations. From those operations, you can generate your model. A really simple application of this is a graph, which is just a set of vertices and edges. So if you have these operations
VertexAdd { id: fooID, label: foo }
VertexAdd { id: barID, label: bar }
VertexAdd { id: quuxID, label: quux }
EdgeAdd { src: fooID, dst: quuxID }
where the ids generated are unique (say, uuids). The corresponding model is just the graph: Graph {
vertices: {foo, bar, quux},
edges: {{foo, bar}},
}
You might imagine operations to remove vertices and edges as well. Many operations are commutative (add vertex), but some pairs of operations don't commute, e.g., add vertex/remove vertex. As I understand that, that's the central problem this paper is trying to address via CRDTs. What I'm saying is that the total ordering imposed by the server makes this concern moot, because the server always picks one sequential ordering.Let's say there are multiple simultaneous clients editing our graph structure above at the same time. Let's consider an example of a pair of operations that commute. One of them decides to remove vertex quux while the other adds an edge that connects quux to foo:
client1:
VertexDelete { id: quuxID }
client2:
EdgeAdd { src: quuxID, dst: fooID }
If both clients attempt to update the server with these new operations, then the server will declare one ordering as true and send the resulting operations back to the clients with the correct ordering. The clients must then adjust. In this case, the outcome is the same regardless of the ordering chosen (because you can't draw an edge to a vertex that doesn't exist).What about the case when two operations are not commutative? Well, that means they are dependent, which in turn means that every client receives them in the order in which they were applied. An add/delete vertex is one such example, because the client doing the delete must acquire a vertex ID to give to the VertexDelete operation, but the only way to acquire the ID is to observer the VertexAdd first.
Another example of dependent operations is attaching meta data to vertices, e.g.,
client1:
VertexAdd { id: anotherID, label: Another }
VertexLabel { id: anotherID, label: Changed }
In this case, VertexAdd/VertexLabel are dependent, but the Label operation requires the vertex ID.Now of course, clients can misbehave. There's no reason why a client couldn't generate the ID and add operations in this order:
client1:
VertexLabel { id: anotherID, label: Changed }
VertexAdd { id: anotherID, label: Another }
But if we can assume well behaving clients, then this should be OK! And even if clients are misbehaving, then the process of turning the operations into a model can simply drop operations like VertexLabel because they reference a vertex that doesn't exist.But this seems no different than "last write wins". The point is that clients want to keep reliably working on their copy without suddenly losing all of their work because of an incorrect/stale assumption that their changes would be accepted. This is typically considered undesirable.
So with non-commutative changes, either you would sync with the server after each change to ensure you can make meaningful progress, or you risk doing lots of meaningless work and having it all undone.
And note that this only applies to the interactions between local operations and server operations. If a client does a bunch of local work and only a small component of it is tied to a previous operation (e.g., one edge in graph), then that doesn't mean the rest of the work gets dropped. That is, if I create a bunch of new vertexes and edges, and only one of those edges connects my new vertex to a vertex already present in the server's operations and that connection is deleted simultaneously by another client, then I don't lose all of my work---I just lose that one edge. The rest of my graph is in tact.
And of course, if this is an end user application (like, say, Google Docs), then the user should be able to view the history and undo what their collaborator did. But I don't see this as any different than editing a document on Google Docs where one collaborator keeps deleting stuff you did.
But this isn't necessarily true. The point is there's more concurrency and collaboration possible than what you describe, if you restrict the operations in specific ways. This is what CRDTs and CRs do.
> But I don't see this as any different than editing a document on Google Docs where one collaborator keeps deleting stuff you did.
It is different. The collaborator actually saw your changes and decided to remove them. This is very different than your changes arbitrarily being lost by the system itself because its merge behaviour is insufficient.
My best guess is that I've described my design poorly and I'm leaving something critical out. I think the only way for me to figure out what I've left out is to hear more of your objection. :-)
1. There are differences between system conflicts and user conflicts.
Let's say I'm collaborating on a Google doc. Both my collaborator and I start writing in the same part of the document. If the concurrent system is designed "poorly" (for example, it has last write wins semantics), then only one of our writing appears. You can store this history in the system and show that your changes were overwritten by their changes, but this can be frustrating if you are both collaborating on the end of the document.
2. Commutative operations give better guarantees.
Let's take for instance an append only DAG. We restrict this DAG to having 2 operations: creating a vertex (and associated edges) between 2 existing vertices and creating an edge between 2 existing vertices. In this case, even if both my collaborator and I add a vertex between 2 existing vertices, the system will converge on a DAG where both vertices are added. You can follow the same argument with adding edges.
This DAG is an example of a CRDT. It has a limited set of operations (addEdge and addBetween) which allow convergent semantics through concurrent operations
Happy to chat more about this stuff and hope my mobile response makes sense heh.
W.r.t. to (2), yeah, sure, but if clients need to be able to remove things then you lose commutativity and you end up missing the requirements of the system itself. Or is there some other component of a CRDT system that's supposed to handle this? (My big picture takeaway from the OP paper was that it was specifically trying to handle the non-commutative operations.)
Removes are also tricky because if you're just transmitting the remove operation itself, you have to ensure that updates appear to the client in order, or else removes and adds could conflict, causing inconsistent state.
But yeah, this might not work at all for sequences. And of course, there are other issues with respect to the memory requirements of the persistent data structure in the client. Some sort of compaction is probably necessary.
The total order of the central server can make the system simpler and more efficient, but by itself it doesn't solve the problem that Wave has, which is allowing a client to edit his text/message without being interrupted by network latency/interruptions -- imagine typing a letter and having to wait for the server to acknowledge that keypress with >100ms latencies. To solve this problem, you still need some form of xform/merge algorithm that OT and CRDT systems provide.
EDIT: I assumed you were not familiar with OT systems since you didn't mention it in your post, but now that I followed your link I can see that you are. In that light, it seems your comment is more a question about what the tradeoffs are between OT and CRDT systems rather than whether a central server can solve all problems without xform/merge logic.
One tradeoff that comes to mind when thinking about OT and CRDT systems is in the way operations track locations in the datastructure. In OT systems you have offsets (small), in CRDTs you have uuids (large) or dynamically growing identifiers (usually small but possibly large). This has implications for the byte-size of operations or the in-memory datastructure.
Another is that CRDTs have a pruning problem. It has been some time since I looked at CRDTs, but I remember that Wave-style OT didn't have the same problem due to the central server. The pruning problem can cause a CRDT to grow larger than it needs to by forcing it to keep more historic data around just in case it gets an old operation it hasn't seen yet. The central server solves this problem by guaranteeing that it will have sent you all old operations before sending you a newer one. If you know all actors in a n-n system you can also solve this issue, but in an unbounded n-n system I didn't see any way this issue can be solved when I was researching it.
EDIT2: Just want to add that there are lots of other problems that are more practical than theoretical. For example, authorization, authorative copy of the data, REST API, things like that, but that would depend more on the exact use case.
In an unbounded n-n system I still don't see a solution.
Systems with a central server avoid problems of partition tolerance and consistency just fine. It's the availability that's a problem. The system cannot operate without the central server. Whether this is a problem in your use case depends on your use case. Nationwide financial system? Aerospace control system? You probably care a lot about availability and can't use the central server model.
(Although ironically a large part of how stock exchange clearing works is imposing an order on transactions and reconciling the result at the end of the day)
"Clearing" in general is a pre-digital technique for resolving concurrency. Someone writes a cheque for payment and puts a date on it, then exchanges it for goods; at some point later, it reaches the single, central ledger of that customer's bank account which is updated to reflect the transaction.
I'll also reccommend https://martinfowler.com/articles/lmax.html , because it's an interesting architecture with a similar approach of just routing all the operations through in series, but with a number of producer/consumers attached to a ring buffer.
At the limit, this means not having to think about the consistency of your data in terms of the network.
I guess my question is: why is this hard to do? I have another comment[1] in this thread that gives a bit more of a detailed example.
Others in this thread also seem to be implying that my stated design is similar to how Google Wave worked. From Raph's article that I linked my initial comment:
> Almost all practical implementations of OT forego TP2, and solve the problem by limiting concurrency in some way, generally requiring a central server to decide on a canonical ordering of the operations (but still using transformations to let clients apply operations out-of-order just a little bit).
Maybe the parenthetical bit about transformations is the part I'm leaving out that I need to make more explicit? Not sure.
Also, it should be noted that Google's OT does not actually let you work offline for very long. Hence, "apply operations out-of-order just a little bit".
As for your example:
There are two things that algorithms for eventual consistency have to deal with. First, causality. This is your add/remove vertex example. Obviously, a vertex can't be removed before it's added. Any "remove" has to causally follow the "add". Causality, however, is usually dealt with on the transport layer via causal delivery, i.e. it's not really part of the CRDT. The "hard" CRDT/OT case is a set of concurrent operations without any causal order.
In terms of your example, you picked a data structure which is essentially a pair of sets — an easy case. Consider instead an array. Array operations that are not commutative are not necessarily dependent, such as three peers simultaneously adding an item to the same index. There is no causal "order in which they were applied". CRDTs allow you to tiebreak this case, often by deterministically ordering concurrent operations via their causal timestamp and owner UUID. These bits of info are not really part of the data, but they're still necessary to derive a total order. Thus, they tend to be included as part of each CRDT op, ballooning its size. I believe what the paper proposes is separating out this metadata from the operations, then clearing it out to save space when it's no longer necessary. (That is, when all clients have merged past the spot of contention.)
And yes, a central server could perform this tiebreaking step. But in the meantime, the peers talking amongst themselves could be building on their locally-resolved conflicts and find themselves rudely interrupted when the server decides that no, the "truth" they've been assuming for the past minute/hour/day/month is actually the wrong one. (Remember, our goal is arbitrarily long periods of local/P2P/offline communication.)
1. In my world, the clients don't communicate with each
other. Everything has to bounce off the central server
first.
2. And yeah, a graph is the easy case. I definitely agree
that my design doesn't address sequences. If your application
doesn't require sequences, though, then perhaps it's
possible to get away with this simpler solution?2. Sure. Actually, if your design happens to be commutative for every operation after causal delivery, then you've invented what the paper calls a "pure op-based CRDT"! That means it'll work with a server or even through P2P (again, w/causal delivery) just by sending the operations around — no extra information required. I believe your vertex/edge example qualifies as long as each new object is created with a unique ID. However, if there's even a remote chance that two objects could be created with the same ID, then things fall apart, since you could have three concurrent operations of DEL-A, INS-A (the new one), and DEL-A, assuming a starting set that includes A. Depending on the order, you could end up with the set including or excluding A. This is why there are multiple kinds of CRDT sets and why they're not "pure op-based", e.g. Add-Wins and Remove-Wins as mentioned in the paper. I think you'll find that most data structures that are remotely interesting can't really be simplified down to the "pure op-based CRDT".
PS, I made a mistake in the last comment, at least in terms of intent: the server won't be able to tiebreak that concurrency case since the ops have already been applied locally to all the clients and they don't commute. Instead, the server will be forced to either transform the operations (as in Operational Transformation) and then send them off (thus making them de facto commute), or alternatively panic, rebase (replaying part of its history), and force all clients to resync. That is, assuming the clients don't keep around the operation history themselves; if they do, they could perform the rebase/replay locally and the operations would (technically) commute. And this would even work for arbitrary operation-based data structures, including sequences or anything else. But then you have to make sure that replaying history has reasonable complexity, e.g. not O(N^2), or you'll be waiting hours if a concurrent commit resolves at the beginning of your operation stack. Ugh, this stuff is confusing...
Aye, that is indeed the case. (Technically, clients could misbehave/be buggy and insert whatever ID they like, but I think I can live with that. Otherwise, ID generation is handled in a way that guarantees uniqueness.)
> That is, assuming the clients don't keep around the operation history themselves; if they do, they could perform the rebase/replay locally and the operations would (technically) commute. And this would even work for arbitrary operation-based data structures, including sequences or anything else. But then you have to make sure that replaying history has reasonable complexity, e.g. not O(N^2), or you'll be waiting hours if a concurrent commit resolves at the beginning of your operation stack.
I came to this same conclusion as well, and this is exactly what the system does. (Notice that I've slyly gone from "here's this hypothetical design" to "I've actually built something similar to this design." :P) Persistent data structures help a lot here.
I do expect to run into scaling problems with this design, but I think there are many buttons to push (of different kinds) to help with that. In particular, while it is not a document sharing service like Google Docs, it does have the same benefit where each operation belongs to one journal, and one journal makes up one graph, but there are many graphs in the system editable by users. Regardless, I'm really happy I provoked this discussion because I learned a lot!
First, little is said about performance. As the paper explains, the meat of each CRDT is pushed into the eval function, which to my understanding is simply a function over the PO-Log. However, is it always possible to adapt a convergent data structure in such a way that eval takes a reasonable amount of time? I notice that sequences — perhaps the most important data type in CRDT-land! — have not been implemented using this approach. If we assume that a sequence can be retrieved from an insert/delete PO-Log by simply sorting it and removing the deleted operations, does that mean that every eval is O(Nlog(N)) at best? And if your solution is to cache the output sequence as an array, a) how do you ensure a correct mapping between the PO-Log and cache on every new operation, and b) what happens if you lose your data and need to replay your PO-Log from scratch? Can the cache be reconstructed O(Nlog(N)) at worst? A complete guess, but maybe having CRDT-specific bits in the prepare/effect steps is what actually allows CRDTs to be performant in the first place! PO-Log representation is alluringly flexible but seems to come with some hefty tradeoffs.
Second, one of the cleanup steps relies on causal stability, i.e. knowing that each client is ahead of a potentially concurrent op. This is a problem in pure, decentralized P2P environments. First, depending on your network architecture, it's not necessarily possible to identify each peer until they actually start sending messages around. Maybe they got their hands on an early revision of the data and have been chipping away for weeks before going online. Second, nothing prevents a peer from connecting for a bit and then leaving forever, thus ensuring that their last edits will never be causally stable. This can be solved with some centralized logic, but then what's the point of using a CRDT at all?
Finally, and less critically, the lossy cleanup steps make it impossible to retrieve old revisions or identify the author of a particular change.