HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

The next HighLoad++ conference will take place on April 6 and 7, 2020 in Saint Petersburg.
Details and tickets at the link. HighLoad++ Siberia 2019. Hall "Krasnoyarsk". June 25, 12:00. Theses and presentation.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

Sometimes, practical requirements conflict with theory, where important aspects for a commercial product are not accounted for. This presentation outlines the process of selecting and combining various approaches to creating components of Causal consistency based on academic research while considering the requirements of a commercial product. Listeners will learn about existing theoretical approaches to logical clocks, dependency tracking, system security, clock synchronization, and why MongoDB opted for certain solutions.

Mikhail Tyulenev (hereafter – MT): – I will talk about Causal consistency – a feature we worked on at MongoDB. I work in the distributed systems group, and we developed it about two years ago.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

In the process, I had to delve into a considerable amount of academic research because this feature is quite well-studied. It turned out that no single paper fits what is required in production, a database due to the very specific demands that exist in almost any production application.

I will discuss how we, as consumers of academic research, prepare something that we can then present to our users as a finished product that is convenient and safe to use.

Causal consistency. Let's clarify the concepts.

First, I want to briefly explain what Causal consistency is. There are two characters – Leonard and Penny (from the series "The Big Bang Theory"):

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

Suppose Penny is in Europe, and Leonard wants to surprise her with some kind of gathering. And he thinks of no better idea than to remove her from the friend list, send all his friends an update on the feed: "Let's surprise Penny!" (she's in Europe, asleep, cannot see any of this, and cannot notice because she is not there). Eventually, he deletes this post, erases it from the feed, and restores access so that she doesn't notice anything and there is no scandal.
This is all wonderful, but let's assume that the system is distributed, and things went a bit differently. For instance, it might happen that the access restriction on Penny occurred after this post appeared, if the events are not causally linked. Essentially, this is an example of when causal consistency is needed to perform a business function (in this case).

In fact, these are quite non-trivial properties of databases—very few actually support them. Let's move on to the models.

Consistency Models

What is a consistency model in databases anyway? It refers to certain guarantees that a distributed system gives regarding which data and in what sequence the client can receive.

Essentially, all consistency models boil down to how much a distributed system resembles one that operates, for example, on a single node on a laptop. The question is how similar a system operating across thousands of geographically distributed nodes is to a laptop where all these properties are essentially fulfilled automatically.

Therefore, consistency models only apply to distributed systems. All systems that existed before and operated on vertical scaling did not experience such problems. There was a single Buffer Cache, and everything was always read from it.

Strong Model

Essentially, the very first model is the Strong model (or rise ability, as it is often called). This is a consistency model that guarantees that every change, as soon as confirmation of its occurrence is received, becomes visible to all users of the system.

This creates a global order of all events in the database. This is a very strong consistency property, and it is generally very costly. Nevertheless, it is very well supported. It’s just very expensive and slow—hence it is rarely used. This is called rise ability.

There is another, even stronger property supported in Spanner, known as External Consistency. We will discuss it a bit later.

Causal

The following is Causal, exactly what I was talking about. There are several sublevels between Strong and Causal that I won’t discuss, but they all come down to Causal. This is an important model because it is the strongest of all models, providing the most robust consistency in the presence of a network or partitions.

Causals refer to situations where events are connected by a cause-and-effect relationship. They are often perceived as Read your on rights from the client's perspective. If a client has observed certain values, they cannot see values from the past. They start to notice prefix reads. It all comes down to the same thing.
Causals as a consistency model involve a partial ordering of events on the server, where events from all clients are observed in the same sequence. In this case, it involves Leonard and Penny.

Eventual

The third model is Eventual Consistency. This is what absolutely all distributed systems support; it is the minimal model that makes any sense. It means that when we make some changes to the data, they eventually become consistent.

At that moment, it doesn’t say anything; otherwise, it would turn into External Consistency— that would be a completely different story. Nevertheless, this is a very popular model, the most widespread. By default, all users of distributed systems use Eventual Consistency.

I want to provide some comparative examples:

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

What do these arrows mean?

  • Latency. As the strength of consistency increases, it becomes greater for obvious reasons: more writes must be made, and confirmation is needed from all hosts and nodes participating in the cluster that the data is already there. Accordingly, in Eventual Consistency, the response is the fastest because, as a rule, you can even commit in memory, and that would generally suffice.
  • Availability. If we understand this as the system's ability to respond in the presence of network outages, partitions, or failures—fault tolerance increases as the consistency model decreases, as it is sufficient for one host to be alive and provide some data. Eventual Consistency does not guarantee anything about the data—it could be anything.
  • Anomalies. At the same time, of course, the number of anomalies increases. In Strong Consistency, they should almost not exist, while Eventual Consistency can have any number of them. The question arises: why do people choose Eventual Consistency if it contains anomalies? The answer lies in the fact that Eventual Consistency models are applicable, and anomalies exist, for example, over a short period; there is the possibility to use a master for reading and somewhat read consistent data; often there is a chance to utilize strong consistency models. Practically, this works, and often the number of anomalies is time-limited.

CAP Theorem

When you see the words consistency and availability, what comes to mind? Right – the CAP theorem! I want to dispel a myth… It’s not me – it’s Martin Kleppmann, who wrote a great article, a wonderful book.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

The CAP theorem is a principle formulated in the 2000s, stating that Consistency, Availability, and Partition Tolerance: take any two, and you cannot choose all three. This was a certain principle. It was proven as a theorem a few years later by Gilbert and Lynch. Then it began to be used as a mantra – systems started to be classified into CA, CP, AP, and so on.

This theorem was actually proven for the following cases… First, Availability was not considered as a continuous value from zero to one hundred (0 – the system is 'dead', 100 – it responds quickly; we are used to viewing it this way), but as a property of the algorithm that guarantees that in all its executions it returns data.

There is not a word about response time! There is an algorithm that returns data after 100 years – a perfectly available algorithm, which is part of the CAP theorem.
Second: the theorem was proven for changes in the values of the same key, while these changes are a resizable line. This means that in reality they are hardly used because there are other models like Eventual Consistency and Strong Consistency (possibly).

What’s the point of all this? That the CAP theorem, in the form it was proven, is practically inapplicable and rarely used. In its theoretical form, it somehow limits everything. It turns out to be a principle that is intuitively correct but not really proven.

Causal consistency is the strongest model.

What is happening now allows us to achieve all three things: Consistency, Availability through Partitions. In particular, Causal consistency is the strongest consistency model that still works in the presence of Partitions (network disruptions). That's why it represents such a significant interest and why we've engaged with it.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

Firstly, it simplifies the work of application developers. In particular, there is considerable backing from the server: when all records occurring within a single client are guaranteed to come in such an order to another client. Secondly, it withstands partitions.

The inner workings of MongoDB

Remembering that it’s lunchtime, we move into the kitchen. I’ll explain the system model, specifically what MongoDB is for those hearing about such a database for the first time.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

MongoDB (hereinafter referred to as 'MongoDB') is a distributed system that supports horizontal scaling, meaning sharding; and within each shard, it also supports data redundancy, that is, replication.

Sharding in MongoDB (a non-relational DB) performs automatic balancing, which means each document collection (or 'table' in relational database terms) is divided into chunks, and the server automatically moves them between shards.

The Query Router, which distributes requests, acts as a client for the user, through which they operate. It already knows where and what data resides and directs all requests to the correct shard.

Another important point: MongoDB is single master. There is one Primary – it can take entries that support the keys it contains. Multi-master writes are not possible.

We released version 4.2 – new interesting features were added. In particular, Lucene has been integrated – search – specifically executable Java directly in Mongo, and it has become possible to perform searches through Lucene, just like in Elastic.

And we created a new product – Charts, which is also available on Atlas (Mongo’s own Cloud). They have a Free Tier – you can play around with it. I really liked Charts – data visualization, very intuitive.

Ingredients of Causal consistency

I counted about 230 articles that have been published on this topic – by Leslie Lamport. I will now share with you some parts of these materials from memory.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

It all started with an article by Leslie Lampert, written in the 1970s. As you can see, research in this area is still ongoing. Currently, Causal consistency is gaining interest due to the development of distributed systems.

Restrictions

What limitations exist? This is actually one of the main points, as the constraints imposed by production systems differ significantly from those found in academic articles. Often, they are quite artificial.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

  • First of all, MongoDB is a single master, as I mentioned (this simplifies things a lot).
  • We believe that the system should support around 10,000 shards. We cannot make architectural decisions that would clearly limit this number.
  • We have a cloud, but we assume that users should still have the option to download binaries, run them on their laptops, and everything should work perfectly.
  • We assume that in research, external clients can do whatever they want: MongoDB is open source. Accordingly, clients can be smart and malicious—they may want to break everything. We assume that Byzantine failures can occur.
  • For external clients outside the perimeter, an important limitation is that if this feature is turned off, no performance degradation should be observed.
  • Another point is quite anti-academic: compatibility of previous and future versions. Old drivers should support new updates, and the database should support old drivers.

In general, all this imposes limitations.

Components of Causal consistency

I will now talk about some components. If we consider Causal consistency in general, we can highlight blocks. We selected from works related to various blocks: Dependency Tracking, clock selection, how these clocks can be synchronized, and how we ensure safety—this is an approximate plan of what I will discuss:

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

Full Dependency Tracking

Why is it needed? To ensure that when data is replicated, each record and each data change contains information about what changes it depends on. The very first and naive change is when each message containing a record also includes information about the previous messages:

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

In this example, the numbers in curly braces are record numbers. Sometimes these records are transmitted entirely with their values, and sometimes only certain versions are shared. The essence is that every change carries information about the previous change (it explicitly contains all this information).

Why did we decide not to use this approach (full tracking)? Obviously, because this approach is impractical: any change in a social network depends on all previous changes in that social network, conveying, say, 'Facebook' or 'VK' in each update. Nevertheless, there is a lot of research on Full Dependency Tracking – it predates social networks and works in certain situations.

Explicit Dependency Tracking

The next one is more limited. Here, the transmission of information is also considered, but only that which is explicitly dependent. Usually, the application determines what depends on what. When data is replicated, responses are only given when previous dependencies have been satisfied, that is, shown. This is the essence of how Causal consistency works.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

It sees that record 5 depends on records 1, 2, 3, and 4 – accordingly, it waits before the client can access the changes made by Penny's access decree, when all previous changes have already been processed in the database.

This also does not suit us because there is still too much information, and it will slow things down. There is another approach…

Lamport Clock

They are very old. Lamport Clock implies that these dependencies are folded into a scalar function, which is called the Lamport Clock.

A scalar function is an abstract number. It is often referred to as logical time. With each event, this counter increases. The counter known to the process at that moment sends each message. It is clear that processes can be desynchronized, having completely different times. Nevertheless, through this message exchange, the system somehow balances the clocks. What happens in this case?

I split that large shard in two to clarify: Friends may reside in one node that contains a piece of the collection, while Feed exists in another node that holds a different piece of that collection. How can they avoid being queued? First, Feed will say, "Replicated," and then – Friends. If the system does not provide any guarantees that Feed will not be displayed until the dependencies of Friends in the Friends collection are also delivered, we will encounter the situation I mentioned.

You can see how the logical time on the Feed counter increases:

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

Thus, the main property of this Lamport Clock and Causal consistency (explained through Lamport Clock) is as follows: if we have events A and B, and event B depends on event A, then it follows that the LogicalTime of Event A is less than the LogicalTime of Event B.

* Sometimes it is also said that A happened before B, meaning A occurred before B – this is a relationship that partially orders the whole set of events that have occurred.

The reverse is not true. This is actually one of the main drawbacks of the Lamport Clock – partial ordering. There is a concept of simultaneous events, meaning events where neither (A happened before B) nor (B happened before A). An example could be the parallel addition of someone as a friend by Leonard (not even by Leonard, but by Sheldon, for instance).
This is the property that is often used when working with Lamport clocks: they look specifically at the function and draw conclusions from that – perhaps these events are dependent. Because in one direction, it is true: if LogicalTime A is less than LogicalTime B, then B cannot have happened before A; and if it is greater, then it possibly can.

Vector clocks

The logical advancement of Lamport clocks is Vector clocks. They differ in that each node here contains its own individual clocks, which are transmitted as a vector.
In this case, you see that the zero index of the vector corresponds to the Feed, and the first index of the vector corresponds to Friends (each of these nodes). And now they will be increasing: the zero index of the Feed increases to – 1, 2, 3:

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

What makes Vector clocks better? They help to understand which events are simultaneous and when they occur on different nodes. This is very important for sharding systems like MongoDB. However, we did not choose this, although it is a great tool, works wonderfully, and would probably suit us...

If we have 10,000 shards, we cannot transmit 10,000 components, even if we compress or invent something else – the payload will still be several times smaller than the total volume of this vector. Therefore, grudgingly, we gave up this approach and switched to another.

Spanner TrueTime. Atomic clocks.

I mentioned that there would be a talk about Spanner. It's a cool tool, straight from the 21st century: atomic clocks, GPS synchronization.

What’s the idea? Spanner is Google's system, which has recently even become available to people (they added SQL to it). Each transaction has a certain timestamp. Since the time is synchronized*, each event can be assigned a specific time – atomic clocks have a waiting time after which another time is guaranteed to 'occur'.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

Thus, by simply recording in the database and waiting for some period of time, event Serializability is automatically guaranteed. They have the strongest Consistency model you can imagine – it is External Consistency.

* This is the main problem with Lamport clocks – they are never synchronized in distributed systems. They can drift apart, and even with NTP, they still do not work very well. Spanner has atomic clocks and synchronization, apparently to the microsecond.

Why didn’t we choose this? We do not assume that our users have built-in atomic clocks. When they do become built-in in every laptop, there will be some super-cool GPS synchronization – then yes… But for now, the best available is Amazon, Base Stations – for enthusiasts… So we used other clocks.

Hybrid clocks.

This is essentially what ticks in MongoDB when ensuring Causal consistency. What are they hybrid about? Hybrid is a scalar value, but it consists of two components:

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

  • The first is the Unix epoch (the number of seconds that have passed since the 'beginning of the computer world').
  • The second is an increment, also a 32-bit unsigned int.

That's essentially it. There is an approach where the part responsible for time is constantly synchronized with the clock; every time an update occurs, this part synchronizes with the clock, so the time is always more or less correct, and the increment allows distinguishing between events that occurred at the same moment in time.

Why is this important for MongoDB? Because it allows for certain backup restores to a specific point in time, meaning the event is indexed by time. This is important when certain events are needed; for a database, events are changes in the database that occurred at specific intervals of time.

The main reason I'll only tell you (please, don’t share this with anyone)! We did this because that's how ordered, indexed data looks in MongoDB's OpLog. OpLog is a data structure that contains all changes in the database: they first go into the OpLog, and then are applied to Storage when it's a replicated data or shard.

That was the main reason. However, there are also practical requirements for database development, which means it should be simple – minimal code, as few broken things as possible that need to be rewritten and tested. The fact that our oplogs turned out to be indexed by hybrid clocks greatly helped and allowed for the right choice. It really paid off and somehow worked like magic right on the first prototype. It was very cool!

Clock synchronization

There are several methods of synchronization described in the scientific literature. I'm referring to synchronization when we have two different shards. If there is a replica set, synchronization is not needed: it's a 'single-master' setup; we have an OpLog where all changes are recorded – in this case, everything is already sequentially ordered in the OpLog itself. But when we have two different shards, time synchronization becomes important. This is where vector clocks help more! However, we don't have them.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

The second option is 'Heartbeats'. You can exchange some signals that occur at each time unit. But Heartbeats are too slow; we cannot ensure latency for our client.

True time is, of course, a wonderful thing. But again, this is probably the future… Although in the 'Atlas' it can already be done; there are fast Amazon time synchronizers. But this won't be available for everyone.

Gossiping is when all messages include timestamps. This is roughly what we use. Every message between nodes, drivers, router data nodes, everything for MongoDB – these are some elements, components of the database that contain flowing clocks. They all have a hybrid time value that is transmitted. 64 bits? That’s possible.

How does all of this work together?

Here, I consider one replica set to make it a bit simpler. There is a Primary and a Secondary. The Secondary replicates and is not always fully synchronized with the Primary.

An insert occurs in the Primary with a certain timestamp. This insert increments the internal counter by 11, if that is the maximum. Otherwise, it will check the clock values and synchronize based on those clocks if the clock values are greater. This allows for time ordering.

After it writes, an important moment occurs. Clocks in MongoDB are incremented only when a record is made in the OpLog. This is the event that changes the state of the system. In all classical articles, an event is considered to be when a message reaches a node: when a message arrives, it means the system has changed its state.

This is related to the fact that during the investigation, it is not entirely clear how this message will be interpreted. We know for certain that if it is not reflected in the 'Oplog', it will not be interpreted at all, and the only change in the system state is the entry in the 'Oplog'. This simplifies things for us: both the model is simplified, and it allows for ordering within a single replica set, along with many other useful aspects.

The value that is already recorded in the 'Oplog' is returned – we know that this value is already present in the 'Oplog', and its time is 12. Now, let's say reading starts from another node (Secondary), and it passes the afterClusterTime in the message itself. It states: 'I need everything that happened at least after 12 or at twelve' (see the figure above).

This is what is known as Causal and consistent (CAT). There is a concept in theory that this is some snapshot of time that is self-consistent. In this case, it can be said that this is the system state that was observed at the moment of time 12.

Right now, there is nothing here because it simulates a situation where the Secondary needs to replicate data from the Primary. It is waiting... And then the data arrives – it returns these values back.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

This is more or less how everything works. Almost.

What does 'almost' mean? Let's assume there is a certain person who read and understood how all this works. They grasp that each time a ClusterTime occurs, it updates the internal logical clocks, and then the next entry increments by one. This function takes 20 lines. Let's say this person transmits the largest possible 64-bit number minus one.

Why 'minus one'? Because the internal clocks will be substituted into this value (obviously, this is the largest possible and greater than the current time), then there will be a write to the 'Oplog', and the clocks are incremented by one more – and it will already be the absolute maximum value (there are just all ones, there’s no further, unsaint ints).

It is clear that after this, the system becomes absolutely unavailable for anything. It can only be unloaded, cleaned – a lot of manual work. Full availability:

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

If this is replicated elsewhere, the entire cluster can simply crash. An absolutely unacceptable situation that anyone could easily and quickly create! Therefore, we considered this aspect as one of the most important. How to prevent it?

Our solution is to sign clusterTime.

This is conveyed in the message (before the blue text). But we also started generating a signature (blue text):

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

The signature is generated by a key that is stored inside the database, within a secure perimeter; it is generated and updated (users don’t see any of this). A hash is generated, and each message is signed upon creation, and validated upon receipt.
People might wonder: 'How much does this slow things down?' I emphasized that it should work quickly, especially in the absence of this feature.

What does it mean to use Causal consistency in this case? It shows the afterClusterTime parameter. Without it, it will simply pass values anyways. Gossiping, starting from version 3.6, always works.

If we keep generating signatures constantly, it will slow down the system even when this feature is absent, which does not align with our approaches and requirements. So what did we do?

Do it quickly!

A fairly simple thing, but an interesting trick — I’ll share it; maybe someone will find it interesting.
We have a hash that stores signed data. All data goes through a cache. The cache doesn’t sign a specific time, but a Range. When a certain value comes in, we generate a Range, mask the last 16 bits, and this value is what we sign:

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

By obtaining such a signature, we speed up the system (conditionally) by 65,000 times. It works perfectly: when we conducted experiments, the time for our sequential updates was really cut down by 10,000 times. Of course, when they are out of sync, it doesn't happen. But in most practical scenarios, this works. The combination of Range signing alongside the signature allowed us to solve the security issue.

It is clear that any service strives to avoid downtime. In our case, we believe that the recent failures helped make quay.io better. We took away several key lessons that we would like to share:

Lessons we learned from this:

  • It's important to read materials, stories, and articles because we have many interesting things to share. When we work on a feature (especially now, when we were working on transactions, etc.), we need to read and understand. This takes time, but it is actually very beneficial because it clarifies where we stand. We haven't really invented anything new – we've just taken the ingredients.

    In general, there is a certain difference in thinking when there is an academic conference (like "Sigmon", for example) – everyone focuses on new ideas. What is the novelty of our algorithm? There isn't much novelty here. The novelty rather lies in how we blended existing approaches together. Therefore, the first thing is to read the classics, starting with Lamport.

  • In production, the requirements are completely different. I'm sure many of you face not "spherical" databases in an abstract vacuum, but real, normal things that have issues with availability, latency, and fault tolerance.
  • The last point is that we had to consider different ideas and combine several completely different articles into one approach, together. The idea of signing, for example, actually came from an article discussing the Paxos protocol, which deals with non-Byzantine failures within the authorization protocol, and for Byzantine failures – outside the authorization protocol... Overall, that’s exactly what we ended up doing.

    There is absolutely nothing new here! But once we mixed it all together... It's like saying the recipe for Olivier salad is nonsense because eggs, mayonnaise, and cucumbers have already been invented... It's pretty much the same story.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

I'll conclude here. Thank you!

Questions

Question from the audience (Q): – Thank you, Mikhail, for the presentation! The topic of time is interesting. You use Gossiping. You mentioned that everyone has their own time, everyone knows their local time. I understood that we have a driver – there can be many clients with drivers, also many query planners, and many shards... So what happens to the system if suddenly we have a discrepancy: someone thinks they are a minute ahead, someone – a minute behind? Where will we end up?

MT: – That's a great question! I was actually going to mention shards. If I understand your question correctly, we have this situation: there is shard 1 and shard 2, and reading occurs from these two shards – they diverge, they do not interact with each other because the time they know is different, especially the time they have in the oplogs.
Suppose shard 1 has made a million records, shard 2 has made none at all, and a request comes to both shards. And the first one has an afterClusterTime greater than a million. In this situation, as I explained, shard 2 will never respond.

Q: – I wanted to know how they synchronize and choose one logical time?

MT: – They synchronize very simply. When a shard receives afterClusterTime and does not find a time in the 'Oplog' – it initiates no approved. This means it manually raises its time to that value. This indicates that it does not have events corresponding to that request. It artificially creates this event and thus becomes Causal Consistent.

Q: – What if other events arrive after this that got lost somewhere in the network?

MT: – The shard is designed such that they will not arrive, as it is single master. Once it has recorded, they will not arrive but will come later. It cannot happen that something got stuck somewhere, then it does a no write, and then those events arrive—violating Causal Consistency. When it performs no write, all events should come next (it will wait for them).

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

Q: – I have a few questions regarding queues. Causal consistency assumes that there is a certain queue of actions that need to be performed. What happens if one packet is lost? For example, the 10th, 11th… the 12th is lost, while all the others are waiting for it to be executed. Suddenly our machine crashes, and we cannot do anything. Is there a maximum queue length that accumulates before execution? What fatal failure occurs when losing any single state? Moreover, if we record that there is some previous state, then we should somehow base ourselves on it? But we didn't base ourselves on it!

MT: – That's also a great question! What do we do? In MongoDB, there is the concept of quorum writes and quorum reads. In what cases can a message be lost? When the write is non-quorum or when the read is non-quorum (there may also be some garbage that gets stuck).
We conducted a significant experimental verification regarding Causal consistency, which showed that when reads and writes are non-quorum, violations of Causal consistency occur. That's exactly what you are saying!

Our advice: use at least quorum reads when implementing Causal consistency. In this way, nothing will be lost, even if the quorum write fails… This is an orthogonal situation: if a user wants to avoid data loss, they need to use quorum writes. Causal consistency does not guarantee durability. Durability is guaranteed by replication and the machinery associated with replication.

Q: – When we create an instance that performs sharding (not master, but slave accordingly), does it rely on the Unix time of its own machine or on the time of the ‘master’; is it synchronized the first time or periodically?

MT: – Let me clarify. A shard (i.e., a horizontal partition) always has a Primary. There can be a ‘master’ in a shard and there can be replicas. But the shard must always support writing, because it has to maintain some domain (the Primary is located in the shard).

Q: – So everything depends solely on the ‘master’? Is the ‘master’ time always used?

MT: – Yes. It can be metaphorically said: clocks tick when a write occurs in the ‘master’, in the ‘OpLog’.

Q: – We have a client that connects, and they need to know nothing about time?

MT: – They don’t need to know anything at all! Speaking of how this works on the client side: for the client, when they want to use Causal consistency, they need to open a session. Everything is in there now: transactions in the session and retrieve rights… A session is an ordering of logical events occurring with the client.

If they open this session and state that they want Causal consistency (if the session supports Causal consistency by default), everything works automatically. The driver remembers this time and increments it when it receives a new message. It remembers the response returned by the previous server that provided the data. The next request will include afterCluster (‘time greater than this’).

The client doesn't need to know anything at all! It's completely opaque to them. If people use these features, what does it allow? Firstly, it is safe to read from secondaries: you can write to Primary and read from geographically replicated secondaries, being confident that it works. At the same time, sessions recorded on Primary can even be passed to Secondary, meaning you can use not just one session, but several.

Q: – The topic of Eventual consistency is closely related to a new layer of Computer Science – CRDT (Conflict-free Replicated Data Types). Have you considered integrating these data types into the database, and what can you say about this?

MT: – Good question! CRDT makes sense for conflicts during writes: in MongoDB, there’s a single master.

Q: – I have a question from the DevOps team. In the real world, there are such Jesuitical situations where Byzantine Failures occur, and malicious actors within the protected perimeter begin to poke into the protocol, sending specially crafted packets?

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

MT: – Malicious actors inside the perimeter are like a Trojan horse! They can do many bad things.

Q: – Clearly, leaving a hole in the server, so to speak, through which a whole zoo of elephants can pass and crash the entire cluster forever… It will take time for manual recovery… This is, to put it mildly, incorrect. On the other hand, here's a curious thought: in real life, do we encounter situations where such internal attacks actually occur?

MT: – Since I don’t often encounter security breaches in real life, I can’t say – maybe they do happen. But speaking from a development philosophy standpoint, we think this way: we have a perimeter that protects the folks who handle security – it’s a lock, a wall; and within the perimeter, you can do anything you want. Of course, there are users who can only view and others who have permission to delete the directory.

Depending on the rights, the damage that users can inflict can range from minimal to significant. Clearly, a user with full rights can do anything they want. A user with limited rights can cause considerably less harm. In particular, they cannot break the system.

Q: – Within the secured perimeter, someone is trying to create unexpected protocols for the server to compromise it, and if lucky, the entire cluster… Is it ever that 'good'?

MT: – I have never heard of such things. It's no secret that you can take down a server this way. Taking it down from within, being an authorized user that can write something in a message… In reality, you can't, because it will still need to be verified. There is a way to disable this authentication for users who don't want it – that's their problem; they essentially broke down the walls and you could stuff an elephant in there to trample everything… Actually, you can dress as a technician, come in, and pull it out!

Q: – Thank you for the report. Sergey (Yandex). In Mongo, there is a constant that limits the number of voting members in a Replica Set, and this constant is 7 (seven). Why is this a constant? Why isn't it some kind of parameter?

MT: – A Replica Set can have as many as 40 nodes. There is always a majority. I don't know which version…

Q: – In a Replica Set, you can run non-voting members, but for voting members, the maximum is 7. How do we handle a shutdown in this case if our Replica Set is spread across 3 data centers? One data center can easily go down, and another machine can drop out.

MT: – That's already a bit beyond the scope of the report. It's a general question. Maybe I can talk about it later.

HighLoad++, Mikhail Tyulenev (MongoDB): Causal consistency: from theory to practice

Play video

A little advertisement 🙂

Thank you for staying with us. Do you enjoy our articles? Want to see more interesting content? Support us by placing an order or recommending us to your friends, cloud VPS for developers starting at $4.99, a unique entry-level server alternative that we have created for you: The whole truth about VPS (KVM) E5-2697 v3 (6 Cores) 10GB DDR4 480GB SSD 1Gbps from $19 or how to properly share a server? (options available with RAID1 and RAID10, up to 24 cores and up to 40GB DDR4).

Dell R730xd at half the price in the Equinix Tier IV data center in Amsterdam? Only with us 2 x Intel TetraDeca-Core Xeon 2x E5-2697v3 2.6GHz 14C 64GB DDR4 4x960GB SSD 1Gbps 100TB starting at $199 in the Netherlands! Dell R420 — 2x E5-2430 2.2GHz 6C 128GB DDR3 2x960GB SSD 1Gbps 100TB — from $99! Read about how To build a corporate-class infrastructure using Dell R730xd E5-2650 v4 servers costing 9000 euros for peanuts?

Source: habr.com

Buy reliable website hosting with DDoS protection, VPS VDS servers 🔥 Buy reliable website hosting with DDoS protection, VPS VDS servers | ProHoster