Marcio Cunha

Implementing Resilience in Distributed Databases with Conflict Resolution via CRDTs

Learn how to keep distributed systems consistent and resilient without global locks by using conflict-free replicated data types in practice.

Marcio Cunha•5 min
Also available in:EspañolPortuguês
Summary
  • Conflict-free replicated data types eliminate the need for global locks across unstable networks.
  • Eventual convergence ensures isolated nodes reach the same state automatically after reconnection.
  • Commutative and associative operations allow message arrival order not to alter the final outcome.
  • Counters and observed-removed sets solve common concurrency scenarios without data loss.
  • Storage complexity and metadata consumption increase, requiring scheduled periodic garbage collection.

The Consistency Challenge in Distributed Networks

When building systems that run across multiple servers scattered around the world, we face an inescapable physical barrier: network latency and the possibility of temporary connection drops. Traditionally, databases use locks to ensure that two people do not modify the same data at the same time, preventing inconsistencies. In practice, this means that if a submarine cable linking a server in Brazil to another in Europe breaks, the entire system must halt or refuse writes to prevent divergence. This rigidity protects data, but destroys availability and the user experience during moments of instability.

To bypass this obstacle, modern architects adopt eventual consistency models, where each server accepts writes independently, even while temporarily isolated from others. The major problem with this approach arises when the network stabilizes and servers need to synchronize their information. If user A changed a product price in São Paulo and user B changed the same price in Tokyo at the same second, which value should prevail? Without a clear mathematical rule, the system collapses or overwrites data arbitrarily, generating severe operational losses for the business.

The Concept and Mathematics Behind CRDTs

Conflict-free Replicated Data Types, known by the acronym CRDT, emerge as an elegant solution to this software engineering dilemma. They are special data structures that can be updated on any replica completely independently and concurrently, without any central coordination or network locking. In practice, this means two servers can receive simultaneous modifications and later combine their states through a mathematical function that guarantees both will arrive at the exact same final result, regardless of the order in which messages arrived.

For this magic to work in computer architecture, the underlying mathematical structure must obey rigid algebraic properties, such as commutativity, associativity, and idempotency. In simple terms, commutativity ensures that the order of factors does not alter the product, meaning if message X arrives before Y on one node and the order inverts on another, the final result remains identical. Idempotency ensures that applying the same update multiple times produces the same effect as applying it just once, protecting the system against duplicate packet delivery in unstable mobile networks.

Practical Implementation of a State-Based Counter

To understand how a CRDT works in everyday code, let us examine the conceptual implementation of a distributed counter of the PN-Counter type, which allows both concurrent increments and decrements. Each network node maintains an internal vector corresponding to the total number of nodes in the system, recording the volume of changes made by each participant in isolation. When we need to query the total counter value, the system simply sums all positions of this vector, obtaining a unified and precise view without querying a centralized authority.

class PNCounter:    def __init__(self, node_id, total_nodes):        self.node_id = node_id        self.P = [0] * total_nodes        self.N = [0] * total_nodes    def increment(self, val=1):        self.P[self.node_id] += val    def decrement(self, val=1):        self.N[self.node_id] += val    def value(self):        return sum(self.P) - sum(self.N)    def merge(self, remote_p, remote_n):        for i in range(len(self.P)):            self.P[i] = max(self.P[i], remote_p[i])            self.N[i] = max(self.N[i], remote_n[i])

The code snippet above demonstrates the simplicity and robustness of the merging operation, known in the distributed ecosystem as the merge method. When two nodes exchange information, the algorithm updates the local vector by always picking the highest value found between the local state and the remote state for each specific position. In practice, this means if a node missed a previous message, it absorbs the partner's more advanced state in a fully deterministic manner, eliminating any chance of concurrency conflicts or silent loss of important updates.

State-Based versus Operation-Based Models

In distributed systems engineering, CRDTs fundamentally divide into two major architectural categories: state-based ones, called CvRDTs, and operation-based ones, known as CmRDTs. State-based models work by transmitting the complete data structure or an integral copy of its current state whenever synchronization occurs between network nodes. In practice, this means communication is simple to implement because if a message gets lost along the way, the next successful update automatically corrects all past accumulated divergences.

On the other hand, operation-based models transmit only the atomic command that generated the modification, such as the exact instruction to add an element to a list or increment a variable by one unit. This approach consumes much less network bandwidth compared to sending entire heavy structures, making it ideal for environments with limited or expensive connections. However, it requires reliable and causal delivery guarantees for messages from underlying transport protocols, because if an insertion command arrives before its respective initialization, the system can corrupt logical state and fail silently.

Common Pitfalls and Hidden Memory Costs

Despite elegantly solving the complex problem of consistency in high-availability environments, CRDTs exact a significant operational price in terms of computational resource consumption. Because structures must store historical metadata and tracking vectors to guarantee correct mathematical convergence, disk space and RAM usage grow proportionally to the number of nodes and update frequency. In practice, this means a system using observed-removed sets can accumulate thousands of logical tombstones of deleted items, requiring complex cleanup routines to prevent total machine resource exhaustion.

Another essential consideration in architectural planning involves modeling application data so it adequately fits the constraints imposed by CRDT mathematical operations. Not every business domain can be easily translated into commutative and associative structures without imposing severe limitations on transactional validation logic. Operations depending on strict real-time uniqueness constraints, such as ensuring two users do not register the same email address simultaneously on distinct servers, remain extremely difficult to implement without resorting to traditional coordination and distributed consensus mechanisms.

Final Considerations for Systems Architects

The adoption of conflict-free replicated data types represents a profound mindset shift in modern software engineering, trading rigid control for mathematical determinism. By accepting that the network is inherently flawed and temporal decoupling is inevitable in large-scale architectures, we manage to build truly resilient applications capable of operating uninterrupted. The success secret lies in carefully evaluating trade-offs of memory consumption and domain complexity before applying this technology, ensuring resilience brings real value to the business without creating unsustainable technical debt in the long run.