Postgres is Crying: A Zero-Downtime, 15TB Migration to BigQuery with Python & Arrow

Postgres is Crying: A Zero-Downtime, 15TB Migration to BigQuery with Python & Arrow

# python# database# cloud# architecture
Postgres is Crying: A Zero-Downtime, 15TB Migration to BigQuery with Python & ArrowJun Matsui

The PagerDuty alert screamed at 2:17 AM. It wasn't a gentle nudge; it was the digital equivalent of a...

The PagerDuty alert screamed at 2:17 AM. It wasn't a gentle nudge; it was the digital equivalent of a fire alarm in a munitions factory. Host high disk I/O wait. I didn't even need to open Grafana. I could feel it. The sickly green line for iowait would be pegged at 99%, our primary PostgreSQL replica would be lagging by 45 minutes, and our analytics team, bless their hearts, would be running a monster GROUP BY on a 3-billion-row table, effectively DDOSing our production database.

PostgreSQL is a masterpiece of engineering. For transactional workloads (OLTP), it's a reliable, acid-washed Swiss Army knife. But we had pushed our 15TB instance past its breaking point. It was now serving as a data warehouse, and it was choking. Queries that once took seconds were now timing out after 30 minutes. The business needed faster insights, and engineering needed to sleep through the night.

The diagnosis was clear: we were using a screwdriver to hammer a nail. We needed a real hammer. We needed a columnar data warehouse. We needed Google BigQuery.

But how do you move a 15-terabyte, fire-breathing dragon from one cage to another while it's still breathing fire, without anyone noticing?

This is that story.

The Core Problem: Row vs. Columnar Storage

Before we dive into the code, let's build a mental model. Why was Postgres struggling? Because it's a row-oriented database.

Imagine your data is a massive, disorganized dresser full of receipts.

  • Postgres (Row-Oriented): To find the total amount spent on "coffee" last year, Postgres has to pull out every single drawer, open every receipt (row), look at all its details (date, vendor, items, amount), and only keep the "amount" if the vendor is "coffee shop". It's incredibly inefficient for analytics.

  • BigQuery (Columnar): BigQuery organizes the dresser differently. It has one drawer just for "vendors," another just for "amounts," and another for "dates." To get the total coffee spend, it just opens the "vendors" drawer to find all the "coffee shop" entries, notes their positions, and then goes to the "amounts" drawer to grab only the amounts at those exact same positions. It doesn't even touch the other data.

This is why a 45-minute Postgres query can become a 650-millisecond BigQuery query. It's a fundamental architectural advantage for analytics.

The Zero-Downtime Game Plan

Our migration strategy had three non-negotiable phases, running in parallel:

  1. The Great Bulk Export: Get the historical 15TB of data out of Postgres and into BigQuery. This is the heavy lift.
  2. The Live CDC Stream: While the bulk export is chugging along, we need to capture every new INSERT, UPDATE, and DELETE happening on the live database. This is our Change Data Capture (CDC) stream.
  3. The Cutover & Reconciliation: Once the bulk load is done, apply the changes from the CDC stream to the new BigQuery tables, run a paranoid data integrity check, and then, with surgical precision, flip the switch.

Here's a high-level view of the data flow:

              +--------------------------------+
              |                                |
              |      PostgreSQL Primary DB     |
              |          (15 TB, on fire)      |
              |                                |
              +----------------+---------------+
                               |
      +------------------------+--------------------------+
      |                                                   |
      | (Phase 1: Bulk Export)                            | (Phase 2: CDC Stream)
      v                                                   v
+-----+------------------+                     +------------------------+
| Python ETL Worker      |                     |   CDC Listener         |
| (Streaming with Arrow) |                     |   (e.g., Debezium,     |
+-----+------------------+                     |    Logical Replication)|
      |                                        +----------+-------------+
      | Chunks of data                                    |
      | as Parquet files                                  | Live changes
      v                                                   v
+-----+------------------+                     +------------------------+
| Google Cloud Storage   |                     |   Pub/Sub or Kafka     |
| (Staging Area)         |                     +------------------------+
+-----+------------------+                                |
      |                                                   |
      | Load Job                                          v
      v                                        +------------------------+
+-----+--------------------------------------+ | Cloud Function/Stream  |
|                                            | | Processor              |
|             Google BigQuery                +-> (Applies changes)    |
|                                            | +------------------------+
|                                            |
+--------------------------------------------+
Enter fullscreen mode Exit fullscreen mode

Phase 1: The Bulk Export - Python & Arrow to the Rescue

Your first instinct might be to use pandas.read_sql and pandas.to_gbq. Don't. You will summon the OOM (Out Of Memory) Killer, and it will show no mercy. Trying to load gigabytes of data into a single DataFrame is a recipe for disaster.

πŸ† Achievement Unlocked: Evaded the OOM Killer

We need to stream. We'll pull data from Postgres in manageable chunks, convert it to a hyper-efficient in-memory format (Apache Arrow), write it to a columnar file format (Parquet), and upload it to Google Cloud Storage (GCS). BigQuery can then load data from Parquet files in GCS with blinding speed.

Here’s the production-grade Python code that does the heavy lifting.

# requirements: psycopg[binary], pyarrow, google-cloud-storage, google-cloud-bigquery
import os
import uuid
from contextlib import contextmanager
from datetime import datetime, timezone

import pyarrow as pa
import pyarrow.parquet as pq
import psycopg
from google.cloud import bigquery, storage

# --- Configuration ---
# Use env variables in production!
PG_CONN_STRING = "postgresql://user:password@host:port/dbname"
GCS_BUCKET_NAME = "your-bq-staging-bucket"
GCS_PREFIX = f"pg_migration/{datetime.now(timezone.utc).strftime('%Y-%m-%d')}"
BQ_DATASET_ID = "your_dataset"
BQ_TABLE_ID = "your_migrated_table"
CHUNK_SIZE = 500_000  # Number of rows to process at a time

# --- Database & Cloud Clients ---
# In a real app, manage these clients more robustly (e.g., singletons)
storage_client = storage.Client()
bq_client = bigquery.Client()
gcs_bucket = storage_client.bucket(GCS_BUCKET_NAME)

@contextmanager
def get_server_side_cursor(table_name: str, columns: list[str]):
    """
    Creates a server-side cursor to stream data from Postgres without
    loading it all into the client's memory. This is CRITICAL.
    """
    conn = psycopg.connect(PG_CONN_STRING, autocommit=True)
    # A server-side cursor is named. The client only ever holds a small buffer.
    cursor_name = f"migration_cursor_{uuid.uuid4().hex}"
    cursor = conn.cursor(name=cursor_name)
    cursor.itersize = CHUNK_SIZE  # How many rows psycopg fetches at once internally

    print(f"INFO: Starting stream from table {table_name}...")
    try:
        column_str = ", ".join(f'"{c}"' for c in columns)
        cursor.execute(f"SELECT {column_str} FROM {table_name}")
        yield cursor
    finally:
        print("INFO: Closing connection.")
        cursor.close()
        conn.close()

def upload_to_gcs(file_path: str, gcs_path: str):
    """Uploads a local file to GCS."""
    blob = gcs_bucket.blob(gcs_path)
    blob.upload_from_filename(file_path)
    print(f"INFO: Uploaded {file_path} to gs://{GCS_BUCKET_NAME}/{gcs_path}")

def stream_pg_to_gcs_parquet(table_name: str, columns: list[str]):
    """
    The main orchestration function.
    - Streams from Postgres using a server-side cursor.
    - Converts chunks to Arrow RecordBatches.
    - Writes to local Parquet files.
    - Uploads to GCS.
    """
    gcs_paths = []
    with get_server_side_cursor(table_name, columns) as cursor:
        chunk_num = 0
        while True:
            # fetchmany is the key to streaming.
            rows = cursor.fetchmany(CHUNK_SIZE)
            if not rows:
                break

            chunk_num += 1
            print(f"INFO: Processing chunk {chunk_num} with {len(rows)} rows...")

            # Zero-copy conversion to an Arrow Table
            # PyArrow can often directly map Postgres memory to Arrow memory.
            try:
                record_batch = pa.RecordBatch.from_pylist(
                    [dict(zip(columns, row)) for row in rows]
                )
                arrow_table = pa.Table.from_batches([record_batch])
            except Exception as e:
                print(f"ERROR: Could not convert chunk to Arrow: {e}")
                # Add dead-letter queue logic here
                continue

            # Write to a local Parquet file
            local_filename = f"/tmp/{table_name}_chunk_{chunk_num}.parquet"
            pq.write_table(arrow_table, local_filename, compression='SNAPPY')

            # Upload to GCS
            gcs_path = f"{GCS_PREFIX}/{table_name}/part-{chunk_num:05d}.parquet"
            upload_to_gcs(local_filename, gcs_path)
            gcs_paths.append(f"gs://{GCS_BUCKET_NAME}/{gcs_path}")

            os.remove(local_filename)

    print("INFO: Stream finished.")
    return gcs_paths

def load_gcs_to_bigquery(gcs_paths: list[str]):
    """Kicks off a BigQuery load job from GCS."""
    job_config = bigquery.LoadJobConfig(
        source_format=bigquery.SourceFormat.PARQUET,
        write_disposition=bigquery.WriteDisposition.WRITE_TRUNCATE, # Overwrite table
    )

    table_ref = f"{bq_client.project}.{BQ_DATASET_ID}.{BQ_TABLE_ID}"
    load_job = bq_client.load_table_from_uri(gcs_paths, table_ref, job_config=job_config)
    print(f"INFO: Starting BigQuery load job {load_job.job_id}")

    load_job.result() # Waits for the job to complete.
    print("INFO: BigQuery load job finished.")

if __name__ == "__main__":
    # In a real scenario, you'd get these columns from information_schema
    # and handle data type mapping carefully!
    table_columns = ["id", "user_id", "event_type", "event_timestamp", "payload"]

    gcs_file_paths = stream_pg_to_gcs_parquet("events", table_columns)

    if gcs_file_paths:
        load_gcs_to_bigquery(gcs_file_paths)
Enter fullscreen mode Exit fullscreen mode

Phase 2 & 3: The CDC Stream & The Terrifying Cutover

While our Python script was happily chunking through terabytes of history, the production database was still taking live writes. We used PostgreSQL's logical replication to stream these changes into a Pub/Sub topic. A simple Cloud Function listened to this topic and applied the changes to our BigQuery table using MERGE statements.

This architecture diagram shows the two paths running in parallel:

Database Migration Architecture and Pipeline Sequence Diagram

πŸ” Click to View & Pan/Zoom in Mermaid Live Interactive Editor β†—

The Moment of Truth: Data Integrity Validation

You cannot skip this step. Hope is not a strategy. Before flipping the DNS, you must mathematically prove that the data is identical.

πŸ”₯ Seniority Check: Do you know why COUNT(*) is not enough? Because UPDATE operations can change data without changing the row count.

We need to checksum the data. For numeric and timestamp columns, a simple SUM() or AVG() can work. For string or complex data, a hash function is better. BigQuery doesn't have BIT_XOR like Postgres, but we can use CHECKSUM_AGG or build our own aggregate hash.

Here's a simple Python validator:

# A simplified validation script
def validate_data_integrity(pg_table, bq_table, column):
    """Compares checksums between Postgres and BigQuery."""

    # --- Postgres Checksum ---
    with psycopg.connect(PG_CONN_STRING) as conn:
        with conn.cursor() as cur:
            # BIT_XOR is a great way to 'fingerprint' a numeric column
            cur.execute(f"SELECT BIT_XOR(CAST({column} AS BIGINT)) FROM {pg_table};")
            pg_checksum = cur.fetchone()[0]
            cur.execute(f"SELECT COUNT(*) FROM {pg_table};")
            pg_count = cur.fetchone()[0]

    # --- BigQuery Checksum ---
    bq_client = bigquery.Client()
    # BigQuery's equivalent for BIT_XOR
    query = f"""
        SELECT BIT_XOR(CAST({column} AS INT64)) as checksum, COUNT(*) as count
        FROM `{bq_table}`
    """
    query_job = bq_client.query(query)
    results = query_job.result()
    for row in results:
        bq_checksum = row.checksum
        bq_count = row.count

    print("--- DATA INTEGRITY CHECK ---")
    print(f"Postgres Count: {pg_count}, BigQuery Count: {bq_count}")
    print(f"Postgres XOR:   {pg_checksum}, BigQuery XOR:   {bq_checksum}")

    if pg_count == bq_count and pg_checksum == bq_checksum:
        print("βœ… SUCCESS: Data is consistent!")
        return True
    else:
        print("❌ FAILURE: Data mismatch detected!")
        return False

# Run this right before the final cutover
# validate_data_integrity("public.events", "your_project.your_dataset.your_migrated_table", "id")
Enter fullscreen mode Exit fullscreen mode

The "Gotchas" The Docs Don't Tell You

This was not a smooth ride. Here are the landmines we hit:

  1. Timezones are the Devil's Work: Postgres TIMESTAMPTZ stores timestamps in UTC internally. BigQuery's TIMESTAMP is also UTC. Sounds great, right? Wrong. If your application layer was accidentally inserting timezone-naive timestamps, Postgres would assume your server's local timezone. BigQuery will not. You must audit every single timestamp column and ensure your ETL pipeline explicitly sets the timezone to UTC.

  2. The Hidden Costs of GCS: Storing 15TB in GCS for a few days isn't too expensive. But remember to set a lifecycle policy to delete the staging Parquet files after the BigQuery load is successful. Otherwise, you'll have a nasty surprise on your next cloud bill.

  3. BigQuery MERGE Limits: BigQuery's MERGE statement is powerful but has DML limits. If your CDC stream is too high-volume, you can hit quota errors. You need to batch your changes and run the MERGE operation every few minutes, not on every single event.

πŸ•ΉοΈ Mini-Quiz: A query that was fast on Postgres is suddenly slow on BigQuery. What's the most likely culprit? (Click to reveal)

> The Culprit: You're likely doing a SELECT * or filtering on a non-clustered/non-partitioned column. In Postgres, if you were fetching a few full rows by their primary key, it was lightning fast. In BigQuery, a SELECT * forces a full scan of every single column, even if you only need two of them. It completely negates the columnar advantage.
>
> The Fix: Be explicit with SELECT column1, column2. And critically, partition your BigQuery table (e.g., by event_timestamp) and cluster it by commonly filtered columns (e.g., user_id). This is the BigQuery equivalent of adding an index, allowing it to prune massive amounts of data from the scan. It's the difference between reading a 1,200-page book line-by-line vs. flipping directly to Chapter 8.

The Payoff: From Pager Alarms to Peaceful Sleep

After the cutover, the difference was night and day.

Metric Before (PostgreSQL) After (BigQuery) Impact
P95 Query Latency 45 minutes (or timeout) 2.1 seconds -99.99%
P50 Query Latency 38 seconds 650 milliseconds -98.2%
Concurrent Analysts ~3 (before locking) 100+ +33x
DBA Heart Rate 140 bpm (on-call) 65 bpm (sleeping) -53.5%
Storage Cost ~$2,500/mo (Managed PG) ~$300/mo (BQ Storage) -88%

The PagerDuty alerts for iowait went silent. Our analytics team was shipping dashboards faster than ever. And for the first time in months, the on-call engineering team could sleep through the night. It was a brutal, complex project, but the payoff was immeasurable.

If you're feeling the heat from a database that's outgrown its purpose, don't just patch the problem. Re-architect for the right tool, arm yourself with streaming patterns, and always, always validate your data.


Interactive Troubleshooting Guide

Facing your own migration woes? Use this decision tree to diagnose the bottleneck.

Migration Troubleshooting Decision Tree Diagram

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚        🚨 INCIDENT: MIGRATION SLOW OR FAILING?         β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
                           β”‚
             β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”΄β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
             β–Ό                           β–Ό
  [Check Source Postgres]     [Check Python Worker]
  β€’ High CPU / Disk I/O?      β€’ High RAM / OOMKilled?
    β”œβ”€ YES ──► Scale replica    β”œβ”€ YES ──► Stream with PyArrow
    └─ NO  ──► Source is fine   └─ NO  ──► Worker is fine
             β”‚                           β”‚
             β–Ό                           β–Ό
     [Check Network]             [Check BigQuery Load]
  β€’ Low throughput?           β€’ Schema validation errors?
    β”œβ”€ YES ──► Snappy Parquet   β”œβ”€ YES ──► Define explicit schema
    └─ NO  ──► Network is fine  └─ NO  ──► Pipeline is optimal!
Enter fullscreen mode Exit fullscreen mode

πŸ” Click to View & Pan/Zoom in Mermaid Live Interactive Editor β†—

Further Reading & Official Docs