Skip to main content

What is a Distributed Database?

What is a distributed database?

A distributed database stores data across multiple remote nodes, rather than in a centralized physical server. Distributed databases are increasingly popular for multi-region operations, scalable workloads, resilient operations, and real-time applications. There are multiple ways to configure a distributed database, with each model offering different benefits.

Why are distributed databases important?

A traditional centralized database stores all its data on one server. A distributed database spreads data and workloads across a network of multiple independent servers, called nodes. To the user, the distributed database’s network is invisible. The user experiences a single, unified database.

You use a distributed database management system (DDBMS) to administer distributed databases. One of a DDBMS’s main functions is to abstract the database’s application layer from the underlying node infrastructure. Doing so allows the DDBMS to manage data routing, query execution, and data consistency across nodes behind the scenes, making sure users have a typical database experience without extra steps.

Organizations adopt distributed database architectures to address several operational requirements:

To manage the scalability demands of modern applications

The hardware capacity of a single server limits the performance of centralized databases. Distributed databases can scale horizontally more readily, meaning they can add entirely new servers on demand. That allows them to handle growing storage and transactional demands.

Resilience against single points of failure

Centralized databases are also limited by the availability of their servers. Hardware failure or network disruptions at a single location take the entire database offline. Distributed databases inherently avoid this issue because they store their data on multiple server nodes, typically in different geographic locations, which are unlikely to be affected by simultaneous hardware, network, or environmental events.

Meet data locality requirements

Globally-accessible applications often benefit by storing copies of data close to regional users to reduce network latency. Latency is the time it takes for a transaction to make a round-trip to and from a database. Many national and international authorities also enforce data sovereignty laws that require specific data to remain within geographic borders.

For example, the European Union’s General Data Protection Regulation (GDPR) mandates how its citizens’ personal data can be used and stored. Administrators can use a distributed database to place data nodes containing such regulated data within the governing jurisdiction, while still maintaining a globally-accessible application.

Read about network latency »

Read about data sovereignty »

Manage real-time and high-throughput workloads

Cloud-based applications that process a high volume of concurrent transactions can perform much better with a distributed database in their backend architecture. The managing DDBMS can help route transactions to one of the most optimal nodes. This approach is common in financial services, e-commerce, and online gaming.

What are the types of distributed databases?

There are several types of distributed databases.

Homogeneous distributed databases

In a homogeneous distributed database, every node in the network stores the same data using the same data model and runs the same DDBMS. This configuration is straightforward to maintain as it is uniform across all nodes. It is a standard choice for many modern, cloud-native deployments.

Heterogeneous distributed databases

The network nodes in a heterogeneous distributed database are not identical. They might run different operating systems, house different data, or most commonly, use different DDBMS software. For instance, you might want to use a heterogeneous distributed database to connect your relational database on one node with a NoSQL document store on another.

Because the underlying databases operate differently, a heterogeneous system uses translation layers to abstract the differences in schema and query languages among database nodes to more easily interact as a unified system.

One of the most common heterogeneous configurations is a federated distributed database. The nodes of a federated database can function independently or together. When you query the system, it determines which node can best respond and routes the query accordingly.

Relational distributed databases

Relational distributed SQL databases implement the familiar tabular structure, SQL querying capabilities, and strict ACID (Atomicity, Consistency, Isolation, Durability) compliance of traditional relational systems. They use horizontal sharding strategies to maintain these features in a decentralized model.

Instead of using a single server to manage a relational database, the distributed system distributes table rows across multiple computer nodes. Each node, or shard, holds a subset of the data. For this to work, each shard maintains the same schema. The DDBMS routes SQL queries to different shards and aggregates the results. You can scale this configuration horizontally as needed, while maintaining a consistent, unified experience for users. Maintaining ACID compliance requires specific protocols and can result in latency across nodes.

Read about relational distributed databases »

Read about ACID »

NoSQL distributed databases

NoSQL databases do not enforce a unified schema on stored data and are easily scalable. NoSQL databases store unstructured and semistructured data in various models, including document stores, key-value stores, wide-column stores, and graph databases.

Most NoSQL distributed databases operate on the principle of eventual consistency. This approach maximizes system availability at the expense of immediate consistency. All nodes will eventually synchronize data written to any one node, but not necessarily right away. NoSQL systems are well-suited for applications where high throughput is the primary requirement rather than an immediately consistent experience for all users—for example, a social media feed.

Read about NoSQL »

NewSQL distributed databases

NewSQL distributed databases are a class of relational database systems designed to combine the strict ACID compliance of traditional relational databases with the horizontal scalability of NoSQL. NewSQL systems use advanced consensus protocols and distributed transaction routing to maintain consistency in enterprise systems with strict transactional recording requirements. Financial trading and order processing systems often fall into this category.

How does a distributed database work?

Distributed databases rely on a specific set of data management techniques to maintain performance across their network of nodes.

Data distribution strategies

Distributed database systems can apply several different distribution strategies depending on specific technical requirements and use cases.

Horizontal partitioning (Sharding)

Horizontal partitioning in a distributed database, or sharding, splits a database table by rows across multiple nodes. They use a partitioning key called a shard key that determines which data goes where. The two primary methods for this are:

  • Range-based sharding: Divides data based on contiguous ranges of the shard key.
  • Hash-based sharding: Applies an algorithm to the shard key to distribute rows more evenly across nodes.

When configuring a sharded database, administrators must manage the risk of hot partitions. A hot partition occurs when a poorly implemented shard key creates a performance bottleneck by funneling a disproportionate amount of traffic to a single node.

Amazon RDS sharding example diagram

Vertical partitioning

Vertical partitioning splits a table by columns across nodes rather than by rows. For example, instead of storing a portion of complete customer records on each node, you might store contact information for all customers on one node and all order histories on another.

Vertical partitioning is less common than sharding, but it is useful when you have subsets of data that users access at different frequencies. Distributing frequently accessed data across many nodes and containing the rest on just a few can improve overall system performance.

Replication

Replication involves maintaining exact copies of data across multiple nodes. Where sharding splits a dataset across multiple nodes, replication duplicates the entire dataset. Data replication is useful when you need high fault tolerance. If one node fails, the system can direct users to the next-most-available replica. Sharding and replication strategies are often combined to provide both better scaling and availability.

Replication models

When a distributed database system applies replication, it must define how the network nodes handle incoming read and write requests. There are three primary models:

Primary-replica

This is also often called the leader-follower model. The DDBMS designates one node as the primary, or leader, which handles all incoming write operations. Those changes then propagate to all secondary, or follower, nodes. The simplicity of this model improves consistency, but having a single primary node can become a performance bottleneck under write-heavy workloads.

Primary-replica diagram

Multi-primary

Also known as a multi-leader setup, this model allows multiple nodes to accept write operations simultaneously. The advantage over the primary-replica model is higher write throughput.

However, because multiple nodes can simultaneously accept conflicting updates to the same dataset, a multi-primary system requires a reliable conflict-resolution mechanism to maintain data integrity. Common mechanisms include Last Write Wins (LWW) or application-level resolution, in which the DDBMS passes the decision for how to resolve the conflict to the application. However, LWW can result in loss of legitimate writes.

Leaderless

No designated primary node exists in a leaderless model. Any node can accept both read and write operations. These systems use quorum-based coordination to maintain consistency. The leaderless model is often called Dynamo-style coordination.

To make sure that a read operation returns the most recent write, the system must be configured so that the number of nodes required for a successful read (R) plus the number of nodes required for a successful write (W) is greater than the total number of replica nodes (N). This is expressed as the formula R + W > N in database management, where quorums overlap by at least one node.

Distributed transactions and consensus

Completing a distributed database transaction often involves transactions with multiple nodes. The system needs to apply these updates consistently, or it can quickly lose integrity. The system needs mechanisms to make sure that either all participating nodes commit the transaction, or none of them do.

Here are four primary methods for managing distributed database transactions:

Two-phase commit (2PC)

This protocol uses a central coordinator process to decide whether a transaction should complete or abort. The first phase is the voting phase. The coordinator process coordinates activity for all the “worker” processes required to route and commit the transaction. Once prepared, it polls them for a yes/no vote as to whether the transaction will complete based on the local conditions each worker process sees.

The second phase is the commit phase. If all processes vote yes, the transaction commits. If any processes vote no, the transaction stops. The coordinator process communicates the results to all worker processes, and they attempt to rollback and recover their local states. 2PC is resilient to many failure conditions, which has led to widespread adoption. However, 2PC can’t easily handle certain failure conditions. 2PC’s main disadvantage is that it is a blocking protocol. If some local processes vote to commit and the coordinator stops, those local nodes are blocked from completing any other transactions until the coordinator successfully rolls back the transaction or an administrator manually clears it.

Three-phase commit (3PC)

To address the blocking issue found in 2PC, the 3PC method adds an intermediary "pre-commit" phase. This extra step after the voting phase allows participant nodes to agree on the transaction's status before the final commit, reducing the risk of blocking if the coordinator is unresponsive. Two and three-phase commit protocols offer strong ACID compliance.

Consensus protocols

Modern distributed databases frequently rely on mathematical consensus algorithms that enable the database to continue operating if a majority of nodes remain available. Two of the most popular algorithms are Paxos and Raft. Systems such as CockroachDB, Amazon Aurora, etcd, and Spanner use these protocols to manage the state of their distributed systems.

Saga pattern

Instead of locking data across multiple nodes simultaneously, a saga pattern breaks a large transaction into a sequence of smaller, local transactions. If any step fails along the way, the system runs a series of compensating transactions to undo the preceding steps and restore the data. This approach is well-suited to distributed systems with microservices.

What is the CAP theorem?

The CAP theorem states that a distributed database can have at most two of the following three properties at any one time:

  • Consistency: All nodes in the system reflect the exact same data at the same time.
  • Availability: Every request receives a successful, non-error response, even if some nodes are unavailable.
  • Partition tolerance: The system continues to function if the distributed network partitions or breaks communication between nodes.

Network partitions are inevitable in real-world distributed networks, which means partition tolerance is an essential capability and not actually a choice. So in practice, the CAP theorem dictates that database architects must choose whether to prioritize consistency (by refusing requests and returning errors until the network is fixed) or prioritize availability (by returning the available data, even if it might be outdated).

What is the PACELC theorem?

PACELC is an extension of the CAP theorem that describes trade-offs across all operations, not just network failures. PACELC dictates how, in the event of a network Partition, a system must choose between Availability and Consistency (the standard CAP theorem). Otherwise, when the network is functioning normally, the system must choose between Latency and Consistency.

If a database requires consistency, it must wait for data to replicate across multiple nodes before confirming a write, which increases latency. If it prioritizes low latency, it must accept that data will be briefly inconsistent across the network.

What are the key characteristics of a distributed database?

An effectively functioning distributed database should exhibit several specific characteristics:

Abstraction

Distributed architecture should be abstracted, with its complexity hidden from users of a distributed database. They shouldn’t need to know where nodes are or how data is stored in the system and replicates across the network. The system still functions as a single, unified database for them.

Autonomy

Each node within the network can operate independently to some degree.

Fault tolerance

The system as a whole will continue to operate and service requests in the event of individual node failures.

Concurrency control

The system manages simultaneous transactions across the network using mechanisms such as two-phase commit and multi-version concurrency control (MVCC).

Query optimization

The DDBMS should employ a query planner capable of accounting for the multiple physical locations of data and the network cost required to retrieve it to run queries efficiently.

What are the benefits of distributed databases?

Distributed databases provide many important benefits.

  • Horizontal scalability: Distributed databases enable organizations to scale “out” horizontally by adding more nodes.
  • High availability: Distributed databases decrease the risk of single points of failure by storing data across multiple nodes.
  • Geo-distribution: Distributed architecture allows administrators to place data physically closer to the end-users, which can reduce the network latency of read and write operations.
  • Data sovereignty: Distributed databases help organizations keep regulated datasets isolated within regions that enforce data sovereignty laws.
  • Throughput: Distributing data across multiple nodes allows you to perform parallel query and write operations.

What are the challenges of distributed databases?

Although beneficial in many use cases, distributed architecture can pose engineering and operational challenges.

Consistency vs. latency trade-offs

Distributing data on nodes in geographically distant regions can force administrators to choose between consistency, waiting for data to synchronize across all nodes, and latency, or getting a quick response. This decision is highlighted in the PACELC model.

Network partitions

Partition tolerance is a necessary capability of modern distributed database systems because communication disruptions are so routine. Engineers must build a distributed architecture specifically to handle data partitioning, rather than treating it as an edge case.

Operational complexity

Performing schema changes and rebalancing data need planning. Upgrading software in distributed environments requires careful coordination.

Transaction overhead

Because the infrastructure is abstracted, all distributed environments carry some degree of overhead to operate. Protocols such as the two-phase commit add latency. Alternatives such as the saga pattern shift the complexity of managing transactions to the application layer.

Monitoring and observability

Tracking system performance and debugging errors across a distributed architecture can be harder than in a single database. Operations teams typically need to implement distributed tracing tools to follow requests.

Data hotspots

Poorly implemented shard keys can create hot partitions, or data hotspots, where one node handles a disproportionate number of transactions. This can create localized performance bottlenecks.

What are some distributed database use cases?

Organizations across various industries use distributed databases when they outgrow the capabilities of a single-node system.

Financial services

Financial institutions have several use cases for distributed databases. They use them for high-throughput, time-sensitive transaction processing. They also use distributed databases for fraud detection.

E-commerce

Online retailers use distributed systems to manage product catalogs, track inventory, and handle order management at scale. The underlying distributed architecture can maintain performance even during high-traffic spikes.

Gaming

Modern multiplayer games use distributed databases to manage session and player data. Distributing multiple instances of this data geographically helps keep latency low for players in many different regions.

IoT (Internet of Things)

IoT networks generate many continuous streams of information. Distributed databases are well-suited to handling the high write throughput from many concurrent devices, such as industrial IoT sensors.

SaaS applications

Software-as-a-Service (SaaS) providers use a distributed architecture to support multi-tenant applications that need region-specific data isolation. This allows SaaS companies to serve a global customer base while complying with data sovereignty requirements in each region.

Media content delivery and personalization

Media platforms use distributed databases to manage read-heavy workloads for movies, television, and music with geographically distributed caches. Placing data closer to end users allows these platforms to serve customized content more efficiently.

Distributed database vs centralized database

Choosing between a distributed and a centralized database depends on your application's scale, geography, and operational requirements.

When to choose a centralized database

Centralized architecture is typically the best choice for applications with predictable workloads that a single server can accommodate. If your user base is localized, your data volume is manageable, and you want to minimize operational complexity and infrastructure costs, a centralized system is effective.

When to choose a distributed database

Distributed architecture is a better choice when you need to guarantee high availability, process massive volumes of concurrent transactions across multiple regions, or follow strict geographic data sovereignty laws.

How can AWS support your distributed database requirements?

AWS offers a range of database services that allow you to create a distributed architecture for maximum performance and availability:

  • Amazon Aurora is a relational database service that is fully compatible with MySQL and PostgreSQL, allowing existing applications and tools to run without requiring modification. Amazon Aurora is designed to deliver up to five times the throughput of MySQL and three times the throughput of PostgreSQL. Aurora supports cross-Region read replicas.
  • Amazon DynamoDB is a serverless, NoSQL, fully managed database with single-digit millisecond performance at any scale. A DynamoDB global table is comprised of multiple replica tables. Each replica table exists in a different Region, but all replicas share the same primary key schema. When data is written to any replica table, DynamoDB automatically replicates that data to all other replica tables in the global table.
  • Amazon Keyspaces (for Apache Cassandra) is a scalable, highly available, and managed Apache Cassandra-compatible database service. Amazon Keyspaces is serverless and supports applications that require virtually unlimited throughput and storage.
  • Amazon RDS is an easy-to-manage relational database service optimized for total cost of ownership. Amazon RDS Multi-AZ deployments provide enhanced availability and durability for database instances with an SLA of up to 99.95%, making them a natural fit for production database workloads.

Get started with distributed databases on AWS by creating a free account today.

Browse all cloud computing concepts

Browse all cloud computing concepts content here:

Loading
Loading
Loading
Loading
Loading

Did you find what you were looking for today?

Let us know so we can improve the quality of the content on our pages