Horizontally Scaling a Rails Backend with Vitess
An engineering breakdown of distributed consensus, latency isolation, and Raft failover topologies.
Horizontally Scaling a Rails Backend with Vitess
How Shopify evolved the database architecture behind its Shop app to support growing user demand, simplify schema migrations, and distribute MySQL workloads across multiple shards.
How do you scale a production database beyond the limits of a single MySQL instance without rewriting an entire application, introducing unnecessary downtime, or sacrificing data consistency?
Introduction
As the Shop app grew, its Ruby on Rails backend began facing the challenges that come with operating a large, data-intensive application. The team had already improved background jobs, caching, messaging, and database query performance, but the primary MySQL database was becoming a bottleneck.
Large tables consumed multiple terabytes of storage, schema migrations could take weeks, and background jobs sometimes needed to be throttled when database resources were under pressure.
Rather than continuing to make incremental optimizations, Shopify introduced Vitess, an open-source database clustering system built around MySQL.
The migration involved more than splitting tables across servers. It required changes to the data model, query validation, connection routing, transaction boundaries, and deployment strategy.
1. Understanding the Scaling Problem
The initial architecture used Ruby on Rails with a MySQL-based primary datastore. As the application grew, database capacity and operational complexity became increasingly important constraints.
flowchart TD
A[Growing Shop App Traffic] --> B[Rails Backend]
B --> C[Background Jobs]
B --> D[Cache Layer]
B --> E[Message Bus]
B --> F[Primary MySQL Database]
F --> G[Slow Queries and Connection Pressure]
F --> H[Long Schema Migrations]
F --> I[Database Capacity Limits]
The team had already addressed several application-level bottlenecks. However, increasing application capacity alone could not remove the limitations of a primary database that had reached its practical scaling limits.
Why a single database becomes difficult to scale
- Storage growth: Large tables increase the operational cost of backups, maintenance, and schema changes.
- Connection pressure: Too many concurrent database operations can exhaust available connections or server resources.
- Migration duration: Changes to very large tables can take an impractical amount of time.
- Cross-database complexity: Splitting tables across independent databases makes joins and transactions more difficult.
- Background job throttling: Database pressure can reduce the rate at which asynchronous workloads are processed.
Scaling an application server horizontally does not automatically scale its database. If every application instance depends on the same overloaded database, the database remains a shared bottleneck.
2. Evaluating Database Scaling Strategies
Before introducing Vitess, it is useful to understand the architectural alternatives.
Option A: Multi-tenant or pod-based architecture
Separate groups of tenants or workloads into isolated infrastructure units, each with its own database resources.
This approach limits the impact of failures and resource contention between groups. However, the Shop app had a different usage profile from Shopify's merchant-oriented platform, so its architecture required a different approach.
Option B: Move frequently updated data to key-value stores
Key-value databases can distribute certain high-volume workloads efficiently.
The trade-off is that applications may lose some of the flexible indexing and relational query capabilities offered by SQL databases. Data durability and consistency requirements must also be considered.
Option C: Database federation
Distribute different tables or business domains across multiple independent MySQL databases.
flowchart TD
A[Rails Application] --> B[User Database]
A --> C[Orders Database]
A --> D[Configuration Database]
B --> E[Cross-Database Coordination]
C --> E
D --> E
Federation can postpone the limits of a single database, but the application must understand where its data lives. Cross-database joins, transactions, and coordinated schema migrations become more complicated.
Option D: Horizontally scalable database infrastructure
Use a system that manages multiple MySQL instances while providing an abstraction layer for query routing, connection management, and sharding.
This was the approach Shopify selected with Vitess.
3. Why Vitess?
Vitess is an open-source database clustering system that adds a management and routing layer around MySQL.
It supports horizontal sharding, resharding, connection pooling, coordinated schema migrations, and SQL-compatible database access.
Its purpose is to help applications operate across multiple database shards without requiring every application component to manage the underlying topology independently.
Key Vitess components
| Component | Responsibility |
|---|---|
| VTGate | Routes queries, performs query planning, and coordinates database operations. |
| VTTablet | Manages access to individual MySQL instances and provides database-related services, including connection pooling. |
| Keyspace | Represents a logical database whose tables are stored on one or more shards. |
| Shard | A partition of a keyspace containing a subset of its data. |
| VSchema | Defines how tables, keyspaces, and sharding rules are organized. |
| Vindex | Determines how Vitess locates rows or maps values to shards. |
Query execution flow
flowchart LR
A[Rails Application] --> B[VTGate]
B --> C[VTTablet]
C --> D[MySQL Shard]
D --> C
C --> B
B --> A
The application connects through VTGate rather than managing connections to each shard directly. VTGate uses the schema and routing rules to determine where queries should execute.
The application can continue using SQL-oriented database access while Vitess manages much of the complexity associated with distributing data across MySQL instances. This abstraction does not eliminate the need to design queries and transactions for a sharded environment.
4. Choosing the Sharding Key
A sharding key determines how rows are distributed across database shards.
Shopify selected user_id for its user-owned data because most of the relevant tables were associated with an individual user.
This allowed records belonging to a user to be routed according to a consistent key.
For example:
SELECT *
FROM user_profiles
WHERE user_id = 1042;
When the table is sharded using user_id, Vitess can use the value 1042 to identify the relevant shard instead of unnecessarily querying every shard.
Designing tables around the sharding key
Consider a simplified user-owned data model:
CREATE TABLE user_profiles (
id BIGINT PRIMARY KEY,
user_id BIGINT NOT NULL,
display_name VARCHAR(150),
created_at TIMESTAMP
);
CREATE TABLE user_preferences (
id BIGINT PRIMARY KEY,
user_id BIGINT NOT NULL,
preference_key VARCHAR(100),
preference_value TEXT
);
Both tables include user_id, which provides a consistent basis for routing user-owned data.
This is an illustrative schema, not Shopify's production schema.
Why data modeling matters
Some existing tables stored account_id but did not contain user_id, even though the account belonged to a user.
Before sharding those tables by user, the team had to add the missing column and backfill its values.
Because some tables contained billions of rows, these operations were time-consuming and had to be managed carefully while the database continued serving production traffic.
Changing a data model after tables have grown to billions of rows can require substantial migrations and backfills. Identify ownership relationships and enforce the intended sharding key in new tables as early as possible.
5. Preparing the Data Model
Before sharding, the team needed to ensure that the database schema was compatible with the intended data distribution.
The preparation included:
- Identifying tables containing user-owned data.
- Adding the required
user_idcolumns. - Backfilling existing records.
- Reviewing relationships between accounts and users.
- Identifying queries that depended on the old schema.
- Preparing validation mechanisms before changing production routing.
flowchart TD
A[Audit Existing Tables] --> B[Identify User-Owned Data]
B --> C[Add Missing Sharding Keys]
C --> D[Backfill Existing Records]
D --> E[Validate Data and Queries]
E --> F[Prepare for Vitess Sharding]
The main lesson is that a database migration begins well before the database topology changes. Data ownership and query assumptions must already be understood.
6. Building Query Verifiers
One of the most important parts of the migration was validating how application queries would behave after sharding.
A query that works against a single MySQL database may become inefficient, invalid, or operationally risky when the same tables are distributed across multiple shards.
Shopify built application-level query verifiers to detect incompatible patterns before they could cause problems in production.
Types of query verifiers
1. Missing sharding key
Detect queries against sharded tables that omit the information needed to route them efficiently.
-- Potentially requires broader shard access
SELECT *
FROM user_profiles
WHERE display_name = 'Kishore';
-- Includes the sharding key
SELECT *
FROM user_profiles
WHERE user_id = 1042;
The first query may require searching multiple shards if no suitable routing index or lookup mechanism is available. The second includes the selected sharding key.
2. Cross-database transactions
Identify transactions that write to different databases and therefore cannot rely on the same local transaction boundary.
3. Cross-shard transactions
Detect transactions that modify multiple shards when the application's consistency model requires all related changes to remain within a single shard.
4. Cross-shard writes
Find write operations that affect rows across multiple shards, including bulk updates or deletes that do not specify a routing key.
5. Cross-keyspace queries
Detect joins or other query operations that reference tables in different keyspaces and may not be supported by the intended query pattern.
Example of a cross-shard transaction risk
sequenceDiagram
participant A as Application
participant V as VTGate
participant S1 as Shard 1
participant S2 as Shard 2
A->>V: Begin distributed operation
V->>S1: Write record A
S1-->>V: Commit succeeds
V->>S2: Write record B
S2-->>V: Write fails
V-->>A: Partial failure
Distributed transactions require careful handling of failure and atomicity. Shopify chose to avoid cross-shard transactions in the relevant user-data flows rather than relying on them for ordinary application operations.
Run query checks in development, automated tests, and continuous integration. Catching incompatible queries before deployment is safer than discovering them after traffic has been routed to a sharded database.
7. Phase One: Vitessifying MySQL
The first phase introduced Vitess around the existing MySQL database without immediately distributing its data across multiple shards.
Starting architecture
flowchart LR
A[Rails Application] --> B[ProxySQL]
B --> C[MySQL]
Target architecture
flowchart LR
A[Rails Application] --> B[VTGate]
B --> C[VTTablet]
C --> D[Existing MySQL]
The existing database was exposed as an initial, unsharded Vitess keyspace. VTTablet was deployed alongside the MySQL server, and VTGate was configured to make the keyspace accessible to the application.
The data did not need to be distributed across multiple shards at this stage.
Why this phase mattered
- It introduced the Vitess query path before adding sharding complexity.
- It allowed the team to validate connectivity and operational behavior.
- It provided a foundation for later data distribution.
- It allowed the application to adopt the new connection path incrementally.
Shopify reported that this initial transformation required no downtime and was transparent to the application.
8. Safely Switching Production Connections
Changing database routing in production can introduce significant risk. A direct, all-at-once switch makes it harder to isolate problems if the new infrastructure behaves unexpectedly.
Shopify introduced a dynamic connection switcher in the application layer.
The switcher allowed requests to use either the existing ProxySQL route or the new VTGate route, with traffic shifted gradually.
flowchart TD
A[Production Application] --> B{Connection Switcher}
B --> C[Existing ProxySQL Route]
B --> D[New VTGate Route]
C --> E[MySQL]
D --> F[Vitess VTGate]
F --> G[VTTablet]
G --> E
The rollout strategy
- Start with lower-risk workloads, including background jobs with retry mechanisms.
- Route a controlled portion of traffic through VTGate.
- Observe query behavior, errors, and database performance.
- Resolve issues before increasing the traffic percentage.
- Continue until the new connection path is fully deployed.
- Remove the old ProxySQL route when it is no longer needed.
This approach makes it possible to validate a major infrastructure change while retaining control over the rollout.
9. Phase Two: Splitting Data into Keyspaces
Once the application was using Vitess, the next step was to separate tables into logical groups.
The team organized data into three unsharded keyspaces:
| Keyspace | Purpose |
|---|---|
users | Tables containing user-owned data that would eventually be sharded. |
global | Data that was not owned by a particular user. |
configuration | Configuration data and tables needed for sequence-related functionality in the later migration phase. |
At this stage, the keyspaces were still unsharded. This separated the logical organization of data from the later decision to distribute user data across multiple shards.
Moving tables between keyspaces
Vitess provides a MoveTables workflow to migrate tables between keyspaces.
Before running the workflow in production, the team rehearsed the migration in staging environments to uncover configuration problems and software issues.
flowchart TD
A[Initial Unsharded Keyspace] --> B[Identify Table Ownership]
B --> C[Prepare Destination Keyspaces]
C --> D[Move Tables in Staging]
D --> E[Validate Data and Queries]
E --> F[Move Tables in Production]
F --> G[Verify Application Behavior]
A successful test with a small dataset does not guarantee that a migration will work with production-scale tables. Rehearse migration procedures, test failure scenarios, and use realistic data volumes whenever possible.
10. Phase Three: Sharding the Users Keyspace
After separating the data into logical keyspaces, Shopify introduced sharding for the user-owned data.
The resulting architecture consisted of:
- A sharded
userskeyspace. - An unsharded
globalkeyspace. - An unsharded
configurationkeyspace. - A separate sharded lookup keyspace used for lookup Vindexes.
flowchart TD
A[Rails Application] --> B[VTGate]
B --> C[Users Keyspace]
B --> D[Global Keyspace]
B --> E[Configuration Keyspace]
C --> F[User Shard 1]
C --> G[User Shard 2]
C --> H[User Shard N]
B --> I[Lookup Keyspace]
The user-owned tables were distributed according to their sharding rules, while global and configuration data remained in unsharded keyspaces.
This division made data ownership and routing more explicit, but it also introduced new requirements around identifiers, query routing, and uniqueness.
11. Handling Auto-Incrementing IDs with Sequences
Traditional Rails applications commonly use auto-incrementing integer primary keys.
In a sharded database, independent MySQL instances cannot simply generate IDs without coordination if those IDs must be globally unique.
Vitess sequences provide a mechanism for generating identifiers across shards.
For example:
CREATE TABLE orders (
id BIGINT PRIMARY KEY,
user_id BIGINT NOT NULL,
status VARCHAR(50)
);
The id identifies the order, while user_id can determine where the order is stored.
These fields serve different purposes:
- Primary key: Identifies an individual row.
- Sharding key: Determines where the row belongs.
- Sequence: Provides identifiers suitable for use across a distributed database.
A globally unique ID does not automatically tell Vitess which shard contains a row. Queries still need a routing strategy, such as the primary sharding key or an appropriate lookup Vindex.
12. Understanding Vindexes
A Vindex is a Vitess abstraction that helps map values to database shards or locate rows.
Primary Vindexes
A primary Vindex determines how values of the selected sharding key map to shards.
For example, a user identifier can be mapped to the shard responsible for that user's data.
Lookup Vindexes
A lookup Vindex uses a separate mapping to find the shard associated with a value. This is useful when an application needs to locate a row using a value that is not its primary sharding key.
Lookup Vindexes can support uniqueness-related requirements, depending on the lookup configuration and constraints.
Why Vindexes matter
Without a suitable routing mechanism, queries may need to visit multiple shards. This increases resource consumption and can increase latency.
Understanding Vindexes helps developers design queries that reach the required data efficiently.
13. Understanding VSchemas
A VSchema describes the logical organization of a Vitess keyspace and the rules used to route queries.
For sharded keyspaces, it defines how tables relate to sharding keys and Vindexes.
Conceptually, a VSchema helps Vitess answer questions such as:
- Which keyspace contains a table?
- Which column determines the shard?
- How should a query be routed?
- How can a row be located using a lookup Vindex?
The VSchema is an important part of the application's database architecture because it connects logical SQL tables to the physical distribution of data.
14. Managing Schema Migrations at Scale
One of the motivations for adopting Vitess was the time required to migrate very large MySQL tables.
Traditional schema changes can become difficult to manage when tables contain billions of rows and database capacity is already constrained.
Vitess supports online schema migration strategies, including its native migration mechanism and gh-ost.
Shopify initially encountered issues with the Vitess migration strategy while using Vitess v14. In its testing, migrations could be terminated after being throttled for more than ten minutes, for example because of replication lag. The issue was reported and subsequently fixed. The article describes the native Vitess migration strategy as the production approach used after upgrading to Vitess v15.
What to consider during a migration
- Table size and expected migration duration.
- Replication lag and database load.
- Migration throttling behavior.
- Failure recovery and rollback requirements.
- Compatibility with the deployed Vitess version.
- Monitoring and validation during execution.
The broader lesson is to test the actual migration strategy against representative workloads rather than assuming that an online migration will always complete without intervention.
15. What Changed After the Migration?
Shopify reported several operational improvements after horizontally scaling the primary datastore.
| Area | Reported outcome |
|---|---|
| Schema migrations | Operations that previously took weeks could take hours. |
| Database capacity | The team could add shards to expand capacity. |
| Background jobs | Less throttling caused by database capacity pressure. |
| Query management | Routing and compatibility became more explicit through verification. |
| Operations | Database topology could be managed through Vitess abstractions. |
These results are specific to Shopify's environment and migration. They should not be interpreted as guaranteed performance improvements for every application.
The trade-off was additional architectural complexity. Developers needed to understand concepts such as sharding keys, Vindexes, VSchemas, and the constraints associated with cross-shard operations.
16. Major Lessons from the Migration
Lesson 1: Choose your sharding strategy early
Identify data ownership and the intended sharding key before the database grows substantially.
Adding a missing sharding key and backfilling billions of rows can be a long-running operation.
Lesson 2: Build query verifiers
Automated checks help prevent developers from introducing queries that conflict with the sharded data model.
Pay particular attention to missing routing keys, cross-keyspace queries, and transaction boundaries.
Lesson 3: Rehearse every critical migration
Practice the complete procedure in staging, including failure handling, validation, and recovery.
Use realistic data volumes and verify the behavior of the actual tools and versions that will be deployed.
Lesson 4: Design transactions around data ownership
Keeping related writes within a single shard can simplify transaction semantics and reduce the risk of partial failures across distributed resources.
When cross-shard workflows are necessary, explicitly design their consistency and recovery behavior.
Lesson 5: Roll out infrastructure changes gradually
A configurable connection switcher and staged traffic migration provide greater control than changing all production connections simultaneously.
Lesson 6: Treat sharding as an application architecture change
Sharding affects more than database configuration. It influences queries, data modeling, identifier generation, testing, migrations, and developer workflows.
17. When Should You Consider Vitess?
Vitess is worth evaluating when a MySQL-based application has measurable scaling constraints that cannot be adequately addressed through query optimization, indexing, caching, read replicas, or ordinary database capacity increases.
It may be relevant when:
- A primary database has become a persistent capacity bottleneck.
- Very large tables make schema migrations operationally expensive.
- The application can define meaningful data ownership and sharding keys.
- The team needs to distribute database workloads across multiple MySQL instances.
- The organization can support the additional operational and development complexity.
Vitess is not automatically necessary for every application. A smaller system may benefit more from query optimization, indexing, caching, connection-pool tuning, or a larger database instance.
Conclusion
Shopify's migration illustrates how horizontal database scaling is a gradual architectural process rather than a single infrastructure change.
The team first introduced Vitess around the existing MySQL database, safely switched application connections, separated tables into logical keyspaces, and then sharded user-owned data.
Along the way, query verifiers, data-model preparation, staging rehearsals, sequences, Vindexes, and VSchemas helped make the new architecture manageable.
The most important takeaway is that scaling a database requires both infrastructure changes and application-level discipline. Distributed database systems can remove important capacity constraints, but developers must account for data placement, query routing, transaction boundaries, and failure handling.
For teams operating a growing MySQL-backed application, Vitess offers a path to horizontal scaling while retaining a largely SQL-compatible development model—with the understanding that distributed data introduces new concepts that the team must learn and maintain.
Further Reading
kishore
@kishore
Guest technical contributor to NexusBlog engineering community.