Hello! Welcome back to our series on designing high-load distributed systems.
In our last session, we implemented the critical synchronization logic for our shadow table migration. By setting up PL/pgSQL triggers, we ensured that all new INSERT, UPDATE, and DELETE operations on the orders table are automatically replicated to orders_shadow.
With the live data stream handled, we now face the challenge of migrating the historical data. This lesson addresses the learning outcome: Execute and monitor a data backfill process to populate a shadow table without locking the original table. We will design a robust, resumable script to copy potentially billions of rows from the live table to the shadow table without causing downtime or performance degradation.
1. The Challenge: Why Naive Backfills Cause Outages
On a staging environment, one might be tempted to run a simple SQL statement to backfill the data:
-- DO NOT RUN THIS IN PRODUCTION
INSERT INTO orders_shadow (id, order_data, created_at, provider_id, fraud_score)
SELECT
id,
order_data,
created_at,
order_data ->> 'providerId',
(order_data ->> 'fraudScore')::integer
FROM orders;
On a high-load production table, this single command would be catastrophic. It initiates a single, massive transaction that reads the entire orders table. To understand exactly why this leads to an outage, it's essential to understand PostgreSQL's locking behavior during such operations.
Running a safe database migration using Postgres
The article 'Running a safe database migration using Postgres' by Retool provides an excellent explanation of why large backfills are dangerous. It details the specific lock contention issues that arise.
Please read the section titled 'Data backfills'. Focus on the explanation of how Postgres locks work at the session and transaction level, and why a single-transaction backfill on a large, write-heavy table inevitably leads to blocked application writes.
As the article explains, the core problem is lock contention. The long-running backfill transaction acquires row-level locks (or even stronger locks) and holds them until the entire operation commits or rolls back. For a table with millions of rows, this can take hours. During this time, application requests trying to write to those same rows will be blocked, waiting for the locks to be released. This queue of blocked requests quickly cascades, rendering your application unresponsive.
2. The Solution: Batched and Idempotent Processing
The solution, as highlighted in the Retool article, is to break the monolithic task into thousands of small, independent ones. We will process the data in batches.
This approach has a significant consequence: the overall backfill is no longer atomic. A script failure midway through would leave the table partially populated. Therefore, our backfill process must be:
- Resumable: It must be able to pick up where it left off.
- Idempotent: Running it multiple times (or on the same batch multiple times) must not result in duplicate data or errors.
3. Designing the Backfill Strategy
A robust backfill process involves several distinct steps, from defining the scope to executing the copy and monitoring its impact.
Step 1: Define the Backfill Boundary (The "Completion Point")
Our triggers are handling new data, but what is the exact boundary between "historical" data that needs backfilling and "live" data handled by the triggers? We need to define a "completion point."
Migrate from Postgres using dual-write and backfill
The migration guide from TigerData, 'Migrate from Postgres using dual-write and backfill', introduces a formal concept for this boundary. It's a practical guide that thinks through the realities of distributed writers and data lateness.
Please read section 5, 'Determine the completion point T'. This will explain the concept of a 'consistency range' and how to choose a safe point in time to which you can backfill.
As the guide explains, the completion point, let's call it T_completion, is a timestamp in the recent past. It's the definitive line:
- The backfill script is responsible for all rows where
created_at <= T_completion. - The dual-write triggers are responsible for all rows where
created_at > T_completion(and any updates/deletes to older rows that occur after the migration starts).
Choosing T_completion as, for example, one hour before the backfill script starts, provides a safe buffer for any late-arriving data or replication lag to settle.
Step 2: Implement the Batched Copy Script
We will design a script (e.g., in Python using psycopg2 or a shell script with psql) that iterates through the source table in small chunks based on the primary key.
The core of the script is a loop that does the following:
- Select a batch: Fetch the next
Nrows from theorderstable. A common batch size is between 1,000 and 10,000 rows. - Copy and transform the batch: Insert the selected rows into
orders_shadow, applying the necessary transformations. - Make it idempotent: Use an
ON CONFLICTclause to handle cases where a trigger might have already inserted a row that is also in our current batch. - Track progress: Record the last processed ID to make the script resumable.
- Throttle: Pause briefly between batches to minimize the load on the database.
Here is the key SQL query that would be executed inside the script's loop:
-- Assume :last_processed_id and :batch_size are provided by the script
INSERT INTO orders_shadow (id, order_data, created_at, provider_id, fraud_score)
SELECT
id,
order_data,
created_at,
order_data ->> 'providerId',
(order_data ->> 'fraudScore')::integer
FROM orders
WHERE id > :last_processed_id
ORDER BY id
LIMIT :batch_size
ON CONFLICT (id) DO NOTHING;
Analysis of the Query:
WHERE id > :last_processed_id ORDER BY id LIMIT :batch_size: This is an efficient way to walk through a table using its primary key index. Each iteration is a fast, indexed query.ON CONFLICT (id) DO NOTHING: This is the key to idempotency. If the dual-write trigger has already processed a row (e.g., anUPDATEcame in for a row just as our backfill was about to copy it), thisINSERTwould normally fail with a primary key violation. Instead,DO NOTHINGgracefully skips the row, and we let the trigger's version stand. This prevents the script from crashing and ensures correctness.
Step 3: Monitor the Process
Executing a backfill on a live system is not a "fire and forget" operation. Continuous monitoring is essential to ensure safety.
Progress Monitoring:
The script must provide clear visibility into its progress.
- Logging: The script should log its state after each batch, e.g.,
INFO: Processed batch up to order_id=15480000. 45.2% complete. - State Table: A more robust method is to use a dedicated database table to track progress.
After each batch, the script runsCREATE TABLE migration_status ( script_name VARCHAR(255) PRIMARY KEY, last_processed_id BIGINT, updated_at TIMESTAMPTZ );UPDATE migration_status SET last_processed_id = .... This makes the state persistent and easily queryable.
Performance Monitoring:
You must watch the database's health throughout the backfill.
- Lock Contention: Use the
pg_stat_activityview to look for queries in awaitingstate. If your backfill transactions are causing other application queries to wait on locks, your batch size may be too large or your throttling too aggressive.SELECT wait_event, state, query FROM pg_stat_activity WHERE wait_event IS NOT NULL; - Resource Utilization: Monitor CPU, memory, and I/O on the database server. A well-behaved backfill should cause a noticeable but stable and acceptable increase in load, not sharp spikes that threaten the primary workload.
- Throttling: The script should have a configurable delay (e.g.,
sleep(0.1)) between batches. This is a simple but powerful lever. If you see negative performance impacts, you can pause the script, increase the delay, and resume.
The resource Migrate from Postgres using dual-write and backfill also touches on the practicalities of the copy process. While it suggests using a streaming COPY tool, our batched INSERT ... SELECT approach is often simpler to implement and control for general-purpose backfills, and equally safe when batched correctly.
Conclusion
You now have a complete, production-grade strategy for backfilling historical data during a shadow table migration. This method prioritizes safety and observability, ensuring the live application remains healthy throughout the process.
Key Takeaways:
- Naive, single-transaction backfills are dangerous due to long-held locks that block application writes.
- The correct approach is batched processing, which breaks the migration into thousands of small, short-lived transactions.
- Backfill scripts must be idempotent and resumable. Using
INSERT ... ON CONFLICT DO NOTHINGis a key technique for achieving idempotency against concurrent trigger-based writes. - Continuous monitoring of progress, database load, and lock contention is not optional; it is a core part of the process.
- Throttling via small delays between batches is a crucial tool for managing database load.
Preview of the Next Lesson
The dual-write triggers are running, and the backfill script has successfully populated the shadow table with all historical data. The orders and orders_shadow tables are now, for all practical purposes, in sync. The final step is to make the switch. In our next lesson, we will learn how to perform an atomic cutover to a shadow table using a table rename and validate the migration's success.