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.