Marcio Cunha

Geographically Distributed Fault Tolerant Systems with CRDT Conflict Resolution

Learn how to design software architectures capable of operating across multiple continents without data loss using conflict-free replicated data types.

Marcio Cunha•4 min
Also available in:PortuguêsEspañol
Summary
  • Cross-datacenter data replication requires strict trade-offs within the CAP theorem.
  • Specific mathematical structures eliminate the need for centralized locking mechanisms.
  • Eventual consistency ensures all nodes automatically reach the exact same state.
  • Conflict resolution happens deterministically directly at the data storage layer.
  • Write latency decreases significantly when operations occur locally near users.

The Challenge of Consistency at Global Scale

Building systems that work around the world without freezing requires rethinking how we save information. When a user in Tokyo and another in São Paulo edit the same file at the same time, computers must decide who is right. In practice, this means we cannot rely on a single central server to validate every click, because physical distance creates unavoidable delays limited by the speed of light.

Traditional engineering tries to solve this by creating copies of data in multiple places. However, if the undersea cable connecting Brazil to the United States breaks, both halves of the application keep running independently. When the connection returns, the computers must stitch the pieces back together without losing what each user typed in the meantime. This scenario of unstable networks and partitions is the ultimate stress test for any modern software architecture.

The CAP Theorem and the Choice for Availability

There is a mathematical rule in computing called the CAP Theorem, which states that a distributed system can only guarantee two out of three properties at the same time: consistency, availability, and partition tolerance. Since real-world internet connections fail constantly, partition tolerance is mandatory. This forces us to choose between halting the system during failures or accepting that different parts of the network have temporarily divergent views of the data.

In modern architecture, we almost always choose availability. This means the app keeps opening and accepting new registrations even if the main server is unreachable. The inevitable consequence is that, for a few moments, a client might see an outdated bank balance. The secret of engineering is not preventing this divergence, but ensuring that data converges to the correct value as soon as the network stabilizes.

How Conflict-Free Replicated Data Types Work

To solve this puzzle without locking the system, we use mathematical structures known as CRDTs, or Conflict-Free Replicated Data Types. Think of this like a supermarket shopping cart where multiple people can add items at the same time, and in the end all purchases are combined perfectly, regardless of who put what in first. In practice, the CRDT dictates mathematical rules to merge information automatically and predictably.

There are two main types of these structures: operation-based and state-based. State-based ones send the entire history or the current document summary to other servers, which combine the information using a logical rule called join. Operation-based ones transmit only the performed action, such as adding the letter A at position 5. Both ensure that if all servers receive the same notices, they end up with the exact same final result.

Practical Implementation Example with Counters

To visualize the logic in code, we can examine a simple distributed counter written in JavaScript. It allows different servers to independently increase the value before synchronizing the final states. This model avoids bottlenecks and enables offline operation on mobile devices and remote nodes.

class PNCounter {
constructor(nodeId, totalNodes) {
this.nodeId = nodeId;
this.P = new Array(totalNodes).fill(0);
this.N = new Array(totalNodes).fill(0);
}

increment() {
this.P[this.nodeId]++;
}

decrement() {
this.N[this.nodeId]++;
}

value() {
const sumP = this.P.reduce((a, b) => a + b, 0);
const sumN = this.N.reduce((a, b) => a + b, 0);
return sumP - sumN;
}

merge(remoteP, remoteN) {
for (let i = 0; i < this.P.length; i++) {
this.P[i] = Math.max(this.P[i], remoteP[i]);
this.N[i] = Math.max(this.N[i], remoteN[i]);
}
}
}

In this example, each server has its own slot in the increment and decrement lists. When synchronization occurs, the merge function always takes the highest value recorded by each node. This guarantees that no update is lost, even if messages arrive out of order due to network latency.

State Management and Operational Challenges

Despite solving mathematical synchronization, CRDTs introduce an invisible challenge: memory and disk space consumption. Because they need to keep history or metadata about who changed what to successfully merge later, file sizes tend to grow over time. In practice, engineering teams must implement cleanup routines, known as compaction, to discard old traces that have already been acknowledged across the entire network.

Another critical point is the order of operations in the user interface. Even if mathematics guarantees that data becomes identical on all servers, the human experience can be confusing if deleted text reappears because the change came from a slower server. Therefore, besides using smart database structures, developers must design interfaces that clearly communicate saving and synchronization status in real-time.

Final Considerations

Designing geographically distributed systems requires abandoning the comfort of traditional databases that lock everything to guarantee order. Adopting mathematical structures for autonomous conflict resolution allows the construction of resilient, fast, and truly global applications. Investing in decentralized architectures pays off amply when a business needs to operate without interruptions, regardless of network failures or user distance.