
Jun MatsuiThe 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.
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.
Our migration strategy had three non-negotiable phases, running in parallel:
INSERT, UPDATE, and DELETE happening on the live database. This is our Change Data Capture (CDC) stream.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) |
| | +------------------------+
| |
+--------------------------------------------+
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)
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:
π 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")
This was not a smooth ride. Here are the landmines we hit:
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.
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.
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.
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.
Facing your own migration woes? Use this decision tree to diagnose the bottleneck.
ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
β π¨ 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!
π Click to View & Pan/Zoom in Mermaid Live Interactive Editor β