Estimating the probability of faults causing inconsistent state

How likely are faults leading to bad states in a distributed system? Specifically, I am interested in faults causing the constituent components of the system to enter states that are inconsistent with each other. I have sometimes estimated the likelihood of such faults with a variation of the Drake equation, which was originally invented for estimating the number of active and communicative civilizations in the galaxy.

For example, let’s look at a possible fault in a workflow where microservices use HTTP to communicate and where we want to create a user domain object in the User API. Each microservice manages its own state independently of the others. The user object should have an associated user account in the Identity and Management (IdM) service and a workspace object in a third service; here, the meanings of these objects don’t matter as much as their relations. The diagram below shows the expected happy path for user creation at the top and one possible fault at the bottom:

Sequence diagram showing successful and failed attempts to create a user object together with its associations to a user account and a workspace object. The failed attempt is caused by a server crash and introduces inconsistent state between the services.

The APIs offer upsert-style operations where an operation creates an object if it does not exist already, or otherwise updates it; a key given as an input parameter identifies the object. This provides the idempotency property for the operation, which makes the operation safe to retry.

The fault portrayed at the bottom of the diagram happens when the User API crashes (or shuts down ungracefully) while the API is processing the POST /user request. The crash happens just after sending the request to create the user account to the IdM. The request to the Workspace API was never sent. The states of the User API, IdM, and Workspace API become inconsistent with each other: the IdM service has the user account, but the Workspace API lacks the user’s workspace object (it was never requested), and the User API lacks the user object (it was never committed to the database).

Should we worry about such faults? I think this calls for assessing our tolerance for risk: what are the potential impacts of the faults, and how probable are they?

To get some idea of the probability of any fault that introduces inconsistent state within the system, we can estimate the number of their occurrences during one year of operation:

Nfaults = NrequestsInYear × fwriteOps × fwriteCoord × fwriteCrash × fbadState

Where

fwriteOps
The fraction of all operations entering the system that change the state of one or more services. These are write operations.
fwriteCoord
The fraction of fwriteOps that coordinate with other services to achieve any workflow. Having more different service roles would increase the fraction.
fwriteCrash
The fraction of fwriteCoord where any service crashes while processing the write operation, aborting it. The higher your SLA, the smaller the fraction should be.
fbadState
The fraction of fwriteCrash causing any service to enter a state that is inconsistent with other services in the system.

Further, we can estimate the total number of requests in a year (NrequestsInYear) by assuming a constant rate of incoming requests.

We get about 6.3 fault occurrences if we set

Assuming faults occur independently of each other, the number of faults in a year would follow a Poisson distribution with the expectation of Nfaults. The probability of at least one fault during the year is then

P(≥ 1 fault in a year) = 1 - P(0 faults in a year) = 1 − e−Nfaults ≈ 1 − e−6.3 ≈ 99.8%

Here are the parameters filled in, so you can try changing them:

Nfaults
=
× (60 × 60 × 24 × 365 sec/year) ×
×
×
×
=
P(≥ 1 fault in a year)
=
1 − e−Nfaults
=

Given enough time with a steady rate of incoming requests, even rare faults become likely. I find it surprising that a modest rate of 10 requests/sec is likely enough to introduce inconsistent state within the system in one year of operation. What’s worse is that we might not even notice when it happens.

If the impact of the fault exceeds our tolerance for risk, we need fault-tolerance mechanisms that bring the system back to a consistent state. For example, the client can retry a failed POST /user request until it succeeds, which is safe because the operations are idempotent upserts. Alternatively, services can communicate through events delivered by a log-based message broker, which persists the events and redelivers them after a crash.