208 208.fi

High availability and distributed architecture

Some practical limitations of replicated service


A high availability service is available with a high probability. When a service is distributed to multiple independent servers, an individual server receives fewer requests. This improves not only the throughput or latency but also the fault tolerance. Should one of our name servers slow down or hang, others would guarantee continued service.

In theory, the source of a disruption does not actually matter. By far, the most probable fault of our name servers ought to be the regular update of service containers and the content, which will disable one server at a time for a short duration. Because the function of the updated device will be checked before proceeding to the next one, the update can harm at most one device. If necessary, it can be rolled back to the pre-update state. During this time, other servers will keep the service available. In our case, all this happens automatically in our GitLab based ERP system.

The mathematics of availability

The availability of a service can be described with the probability that the service is available at an arbitrary point of time. If the service is distributed to n devices that can fail independently of each other, such that the probability of a single device is unavailable at a point of time is p, then the probability of the entire service (all devices) being unavailable is pⁿ.

For example, an availability target of 99.9999 % or pⁿ=10⁻⁶ allows a service to be unavailable for at most 86.4 milliseconds per day (86,400 seconds).

A two-node cluster reaches this availability target when the failure probability of a single server is p=10⁻³, that is, the device is unavailable for at most 86.4 seconds per day. The target is realistic if the worst-case failure is a network connectivity issue or a planned daily system update taking at most a minute.

In a three-node cluster, the availability target would be reached if p=10⁻², that is, a server will be unavailable for at most 14 minutes and 24 seconds per day.

How does it look like with more severe disruptions, which could take hours or days to fix? If we assume that at most one or two unlikely faults like that can occur simultaneously, we could simply deploy that amount of independent additional servers.

Larger faults, such as a damaged backbone network connection or a fire, can be corrected by deploying a server at an unaffected data centre. This involves human effort. In such case, the availability of a server would be at most p=10⁻¹ (2.4 hours per day), and meeting the availability target would require deploying at least n=6 servers at independent data centres.

Distributed service, central management

It is easy to distribute a service that is based on rarely changing content. For example, the software and certificates on our name servers are updated twice a week. The data payload can be changed when needed, typically rarely.

In a similar way, it is possible to duplicate a certificate validation service or a static web server. But what about services, whose data content is dynamic?

Replicating changes during use

In an e-mail service, the content (user mailboxes) changes during use. SMTP servers are based on message queues, which are basically designed to be fault-tolerant. Incoming messages could be duplicated on multiple servers even before the sender is allowed to close the SMTP connection. Each message has a unique identifier Message-Id. Typically, mailboxes are personal, and a single mailbox isn’t edited simultaneously on multiple devices. E-mail boxes duplicated on different servers could be kept up to date with each other with the help of locks and log records, i.e. by implementing a simple distributed database system. Management is quite easy if changes to an individual mailbox are allowed from a maximum of one server or terminal at a time.

In an online marketplace or booking system, users view and edit shared content at the same time. If there is a limited number of a viral product or event seats available, users should be able to see the current situation and be able to reserve their choice until the transaction including a possible payment has been confirmed or cancelled. The brain of such an application is often a database system that takes care of transactions and locking.

An application that is built on a centralised database can be transformed to a distributed one by replacing the database management system with a distributed one. A well written application is prepared for error situations that lead to transaction rollback, such as deadlocks or lock wait timeouts.

Distributed systems with multiple writable nodes

It is still rather easy to implement and manage a centralised system where all modifications are submitted to one primary server that replicates the changes to hot standby replicas that can serve read-only requests.

High availability through distributed systems?

It is still rather easy to implement and manage a centralised system where all modifications are submitted to one primary server that replicates the changes to hot standby replicas that can serve read-only requests.

However, it is complex to implement a reliable failover meschanism, promoting a replica to a primary server. What if the primary server remains partially reachable and some users keep using it while others have migrated to a replica that has been promoted to be the new primary server?

A distributed system must never reach a split brain situation where users access disjoint primary servers and assume that everyone shares the same view. If a distributed database management system waits every node to acknowledge every single operation, the throughput is determined by the slowest link. Good performance requires the use of buffering, majority decisions and a quick fault resolution.

A real distributed system, where the service nodes of the cluster may quickly switch roles, requires a distributed algorithm or protocol for leader election, an efficient lock management and a recovery mechanism. When a former leader regains connectivity to the other nodes and finds out that it has been demoted, it must either roll back its local changes that others are missing, or attempt to replicate them after replicating the changes it was missing from the cluster. Leader changes could be frequent, and communication failures are possible during the election.

In a distributed system, events may be ordered in very many different ways, or actually partially ordered, because the different parts of the system who have their own clocks may observe a different ordering of events. Also system upgrades become more complex, because during an upgrade, some nodes will be running a different software version.

Developing robust distributed algorithms requires the use of formal methods, such as correctness proofs and verification via model checking and reachability analysis. To manage the state space explosion, a suitable level of abstraction must be found. One of us has touched some of that in a doctoral dissertation.

At the moment we consider the effort of implementing any form of replication larger than the achievable benefit. We’d rather spend a little more manual effort in case the system fails than run a risk that an automated failover mechanism fails. It is easier to manage disruption when the world is halted. The visible world never stops: the production services are deployed in completely different environments in accordance with our layered data security architecture.

To predict and prevent failures of individual servers, our site reliability engineering audits produces improvement advice.

CC BY 4.0