How to Design a Scalable System with Database Sharding and Read-Write Separation

Design a scalable system by splitting data across multiple database nodes (sharding) to distribute write load, while routing read traffic to replicas (read-write separation) to reduce master contention, typically implemented together using middleware like ShardingSphere.

When single-database architectures hit throughput or storage limits, horizontal scaling becomes mandatory. The JavaGuide repository provides comprehensive guidance on implementing these patterns in Java ecosystems, specifically through the docs/high-performance/read-and-write-separation-and-library-subtable.md document. This article distills those practices into an actionable blueprint for building horizontally scalable data layers.

Why Sharding and Read-Write Separation Matter

The Case for Database Sharding

Sharding addresses three critical bottlenecks:

  • Data volume limits – Single tables exceeding tens of millions of rows degrade scan performance and extend backup windows.
  • Write throughput ceilings – A single primary instance cannot sustain massive concurrent insert/update loads.
  • Workload isolation – Vertical sharding places different business modules (user, order, inventory) on separate physical databases, preventing noisy-neighbor effects.

As detailed in docs/high-performance/read-and-write-separation-and-library-subtable.md, sharding splits a logical table into multiple physical tables across database nodes, enabling linear storage and throughput scaling.

The Case for Read-Write Separation

Most internet services exhibit write-few, read-many characteristics. Read-write separation offloads query traffic to replica nodes:

  • Master focuses on transactions – The primary handles writes and critical reads, while replicas serve analytical queries.
  • Latency reduction – Geographic distribution of read replicas places data closer to users.
  • High availability foundation – Replicas provide failover candidates if the master becomes unavailable.

The JavaGuide repository explains the underlying master-slave replication mechanism—where the master writes binary logs (binlog) that replicas apply—in the same high-performance document.

Architectural Strategies for Database Sharding

Horizontal Sharding (Hash and Range)

Horizontal sharding distributes rows across multiple tables or databases. Two primary algorithms dominate:

  • Hash sharding – Applies a hash function to the sharding key (e.g., user_id % 4) to determine the target node. This ensures uniform distribution but complicates range queries.
  • Range sharding – Allocates contiguous key ranges to specific nodes (e.g., users 1-1,000,000 to Node A). This optimizes range scans but risks hot spots.

The JavaGuide file docs/high-performance/read-and-write-separation-and-library-subtable.md catalogs these algorithms under the section covering common sharding strategies.

Vertical Sharding

Vertical sharding splits tables by column groups, placing different domains on separate databases:

  • User data resides on db_user
  • Order data resides on db_order
  • Product catalog resides on db_product

This reduces inter-table contention and allows independent scaling of different business modules.

Hybrid Approaches

Production systems often combine both strategies: vertical sharding isolates domains, while horizontal sharding distributes data within each domain. This hybrid approach balances workload isolation with massive scalability.

Implementing Read-Write Separation

Master-Slave Replication Basics

The standard MySQL implementation uses asynchronous replication:

  1. Master commits transactions and writes to the binary log (binlog).
  2. I/O threads on replicas pull binlog events.
  3. SQL threads on replicas apply events to local data files.

As noted in docs/high-performance/read-and-write-separation-and-library-subtable.md, this mechanism enables read-write separation but introduces replication lag—a critical consideration for consistency.

Routing Strategies

Middleware handles routing transparently:

  • Automatic routing – INSERT/UPDATE/DELETE statements route to the master; SELECT statements route to replicas based on load-balancing algorithms (round-robin, random, or weight-based).
  • Hint-based routing – For scenarios requiring strong consistency, developers force specific queries to the master using hint managers.

Practical Implementation with ShardingSphere

ShardingSphere (specifically Sharding-JDBC) provides a unified solution for both sharding and read-write separation without code changes to business logic. The JavaGuide repository references this in docs/open-source-project/system-design.md under the ShardingSphere section.

Maven Dependency

Add ShardingSphere-JDBC to your project:

<dependency>
    <groupId>org.apache.shardingsphere</groupId>
    <artifactId>shardingsphere-jdbc-core</artifactId>
    <version>5.4.0</version>
</dependency>

YAML Configuration

Configure both read-write splitting and sharding rules in a single YAML file, as demonstrated in docs/high-performance/read-and-write-separation-and-library-subtable.md:

schemaName: demo_ds

dataSources:
  master_ds:
    url: jdbc:mysql://master:3306/demo?serverTimezone=UTC
    username: root
    password: password
  slave_0_ds:
    url: jdbc:mysql://slave0:3306/demo?serverTimezone=UTC
    username: root
    password: password
  slave_1_ds:
    url: jdbc:mysql://slave1:3306/demo?serverTimezone=UTC
    username: root
    password: password

rules:
  - !READWRITE_SPLITTING
    dataSourceNames: [master_ds, slave_0_ds, slave_1_ds]
    loadBalancerName: round_robin
    staticStrategy:
      writeDataSourceName: master_ds
      readDataSourceNames: [slave_0_ds, slave_1_ds]

  - !SHARDING
    tables:
      t_order:
        actualDataNodes: master_ds.t_order_${0..1}
        tableStrategy:
          standard:
            shardingColumn: order_id
            shardingAlgorithmName: order_id_inline
        keyGenerateStrategy:
          column: order_id
          keyGeneratorName: snowflake
    shardingAlgorithms:
      order_id_inline:
        type: INLINE
        props:
          algorithmExpression: t_order_${order_id % 2}
    keyGenerators:
      snowflake:
        type: SNOWFLAKE

This configuration creates a logical demo_ds datasource that automatically routes writes to master_ds and load-balances reads across slave_0_ds and slave_1_ds, while sharding the t_order table across two physical tables using modulo hashing on order_id.

Forcing Master Route for Strong Consistency

When immediate consistency is required after a write, use HintManager to bypass the replica routing, as shown in the JavaGuide examples:

import org.apache.shardingsphere.api.hint.HintManager;

public class OrderService {
    
    public void createAndVerifyOrder(Order newOrder) {
        HintManager hintManager = HintManager.getInstance();
        hintManager.setMasterRouteOnly();  // Force all queries to master
        
        try {
            orderRepository.save(newOrder);
            Order latest = orderRepository.findById(newOrder.getId());
            // Process latest data with strong consistency
        } finally {
            hintManager.close();  // Restore normal read-write splitting
        }
    }
}

This pattern ensures that read-after-write scenarios retrieve the most recent data directly from the master node, mitigating replication lag issues.

Operational Considerations

Aspect Guidance
Consistency Use HintManager for critical reads; consider semi-synchronous replication if latency tolerance allows.
Re-sharding Hash algorithms provide uniform distribution but complicate range queries. ShardingSphere supports online scaling via configuration updates.
Monitoring Track replica lag using SHOW SLAVE STATUS. ShardingSphere exposes Prometheus metrics for per-shard latency.
Backup & Restore Execute logical dumps per physical shard to maintain manageable backup sizes.
Fail-over Promote replicas using MySQL RESET MASTER or orchestration tools (Orchestrator, MHA). ShardingSphere dynamically updates datasource lists without application restarts.

Summary

  • Database sharding partitions data horizontally (by row) or vertically (by domain) to overcome single-node storage and throughput limits.
  • Read-write separation directs write traffic to a master node and distributes read traffic across replicas, reducing contention and improving latency.
  • ShardingSphere (Sharding-JDBC) unifies both patterns in a single configuration layer, requiring zero changes to existing JDBC-based applications.
  • HintManager provides escape hatches for strong consistency when replication lag is unacceptable.
  • The JavaGuide repository (docs/high-performance/read-and-write-separation-and-library-subtable.md) provides the foundational theory and configuration examples for implementing these patterns in production Java systems.

Frequently Asked Questions

What is the difference between database sharding and read-write separation?

Database sharding splits data across multiple physical databases or tables to distribute storage and write load, whereas read-write separation keeps the same dataset on multiple nodes but directs writes exclusively to a master while routing reads to replicas. Sharding addresses data volume and write throughput, while read-write separation addresses read scalability and master node contention.

When should I force queries to the master node instead of allowing replica reads?

Force master routing via HintManager.setMasterRouteOnly() when your application requires read-after-write consistency—for example, immediately after inserting a new order and then querying its status. Replica nodes lag behind the master due to asynchronous replication, so time-sensitive reads must bypass the read-write splitter to access the most current data.

How does ShardingSphere handle both sharding and read-write separation in one configuration?

ShardingSphere treats read-write splitting as a rule that sits above the physical data sources. In the YAML configuration, you first define the master and replica data sources, then apply a !READWRITE_SPLITTING rule to create a logical data source. You then reference this logical data source in your !SHARDING rule's actualDataNodes. This layering allows a single DataSource bean in your application to handle both horizontal partitioning and automatic read-write routing transparently.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →