<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom" xmlns:dc="http://purl.org/dc/elements/1.1/">
  <channel>
    <title>DEV Community: Jun Matsui</title>
    <description>The latest articles on DEV Community by Jun Matsui (@jun-matsui).</description>
    <link>https://dev.to/jun-matsui</link>
    <image>
      <url>https://media2.dev.to/dynamic/image/width=90,height=90,fit=cover,gravity=auto,format=auto/https:%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Fuser%2Fprofile_image%2F2175532%2F2a74d724-980d-4a11-bb55-029c0a878adf.jpg</url>
      <title>DEV Community: Jun Matsui</title>
      <link>https://dev.to/jun-matsui</link>
    </image>
    <atom:link rel="self" type="application/rss+xml" href="https://dev.to/feed/jun-matsui"/>
    <language>en</language>
    <item>
      <title>Postgres is Crying: A Zero-Downtime, 15TB Migration to BigQuery with Python &amp; Arrow</title>
      <dc:creator>Jun Matsui</dc:creator>
      <pubDate>Thu, 08 Oct 2026 21:07:48 +0000</pubDate>
      <link>https://dev.to/jun-matsui/postgres-is-crying-a-zero-downtime-15tb-migration-to-bigquery-with-python-arrow-57ko</link>
      <guid>https://dev.to/jun-matsui/postgres-is-crying-a-zero-downtime-15tb-migration-to-bigquery-with-python-arrow-57ko</guid>
      <description>&lt;p&gt;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. &lt;code&gt;Host high disk I/O wait&lt;/code&gt;. I didn't even need to open Grafana. I could feel it. The sickly green line for &lt;code&gt;iowait&lt;/code&gt; 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 &lt;code&gt;GROUP BY&lt;/code&gt; on a 3-billion-row table, effectively DDOSing our production database.&lt;/p&gt;

&lt;p&gt;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.&lt;/p&gt;

&lt;p&gt;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.&lt;/p&gt;

&lt;p&gt;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?&lt;/p&gt;

&lt;p&gt;This is that story.&lt;/p&gt;

&lt;h3&gt;
  
  
  The Core Problem: Row vs. Columnar Storage
&lt;/h3&gt;

&lt;p&gt;Before we dive into the code, let's build a mental model. Why was Postgres struggling? Because it's a &lt;strong&gt;row-oriented&lt;/strong&gt; database.&lt;/p&gt;

&lt;p&gt;Imagine your data is a massive, disorganized dresser full of receipts.&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;&lt;p&gt;&lt;strong&gt;Postgres (Row-Oriented):&lt;/strong&gt; To find the total amount spent on "coffee" last year, Postgres has to pull out every single drawer, open every receipt (&lt;code&gt;row&lt;/code&gt;), 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.&lt;/p&gt;&lt;/li&gt;
&lt;li&gt;&lt;p&gt;&lt;strong&gt;BigQuery (Columnar):&lt;/strong&gt; 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.&lt;/p&gt;&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;This is why a 45-minute Postgres query can become a 650-millisecond BigQuery query. It's a fundamental architectural advantage for analytics.&lt;/p&gt;

&lt;h3&gt;
  
  
  The Zero-Downtime Game Plan
&lt;/h3&gt;

&lt;p&gt;Our migration strategy had three non-negotiable phases, running in parallel:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt; &lt;strong&gt;The Great Bulk Export:&lt;/strong&gt; Get the historical 15TB of data out of Postgres and into BigQuery. This is the heavy lift.&lt;/li&gt;
&lt;li&gt; &lt;strong&gt;The Live CDC Stream:&lt;/strong&gt; While the bulk export is chugging along, we need to capture every new &lt;code&gt;INSERT&lt;/code&gt;, &lt;code&gt;UPDATE&lt;/code&gt;, and &lt;code&gt;DELETE&lt;/code&gt; happening on the live database. This is our Change Data Capture (CDC) stream.&lt;/li&gt;
&lt;li&gt; &lt;strong&gt;The Cutover &amp;amp; Reconciliation:&lt;/strong&gt; 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.&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;Here's a high-level view of the data flow:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;              +--------------------------------+
              |                                |
              |      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                +-&amp;gt; (Applies changes)    |
|                                            | +------------------------+
|                                            |
+--------------------------------------------+
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h3&gt;
  
  
  Phase 1: The Bulk Export - Python &amp;amp; Arrow to the Rescue
&lt;/h3&gt;

&lt;p&gt;Your first instinct might be to use &lt;code&gt;pandas.read_sql&lt;/code&gt; and &lt;code&gt;pandas.to_gbq&lt;/code&gt;. &lt;strong&gt;Don't.&lt;/strong&gt; 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.&lt;/p&gt;

&lt;p&gt;🏆 &lt;strong&gt;Achievement Unlocked: Evaded the OOM Killer&lt;/strong&gt;&lt;/p&gt;

&lt;p&gt;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.&lt;/p&gt;

&lt;p&gt;Here’s the production-grade Python code that does the heavy lifting.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="c1"&gt;# requirements: psycopg[binary], pyarrow, google-cloud-storage, google-cloud-bigquery
&lt;/span&gt;&lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;os&lt;/span&gt;
&lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;uuid&lt;/span&gt;
&lt;span class="kn"&gt;from&lt;/span&gt; &lt;span class="n"&gt;contextlib&lt;/span&gt; &lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;contextmanager&lt;/span&gt;
&lt;span class="kn"&gt;from&lt;/span&gt; &lt;span class="n"&gt;datetime&lt;/span&gt; &lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;datetime&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;timezone&lt;/span&gt;

&lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;pyarrow&lt;/span&gt; &lt;span class="k"&gt;as&lt;/span&gt; &lt;span class="n"&gt;pa&lt;/span&gt;
&lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;pyarrow.parquet&lt;/span&gt; &lt;span class="k"&gt;as&lt;/span&gt; &lt;span class="n"&gt;pq&lt;/span&gt;
&lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;psycopg&lt;/span&gt;
&lt;span class="kn"&gt;from&lt;/span&gt; &lt;span class="n"&gt;google.cloud&lt;/span&gt; &lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;bigquery&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;storage&lt;/span&gt;

&lt;span class="c1"&gt;# --- Configuration ---
# Use env variables in production!
&lt;/span&gt;&lt;span class="n"&gt;PG_CONN_STRING&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;postgresql://user:password@host:port/dbname&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
&lt;span class="n"&gt;GCS_BUCKET_NAME&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;your-bq-staging-bucket&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
&lt;span class="n"&gt;GCS_PREFIX&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;pg_migration/&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;datetime&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;now&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;timezone&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;utc&lt;/span&gt;&lt;span class="p"&gt;).&lt;/span&gt;&lt;span class="nf"&gt;strftime&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;'&lt;/span&gt;&lt;span class="s"&gt;%Y-%m-%d&lt;/span&gt;&lt;span class="sh"&gt;'&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
&lt;span class="n"&gt;BQ_DATASET_ID&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;your_dataset&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
&lt;span class="n"&gt;BQ_TABLE_ID&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;your_migrated_table&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
&lt;span class="n"&gt;CHUNK_SIZE&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="mi"&gt;500_000&lt;/span&gt;  &lt;span class="c1"&gt;# Number of rows to process at a time
&lt;/span&gt;
&lt;span class="c1"&gt;# --- Database &amp;amp; Cloud Clients ---
# In a real app, manage these clients more robustly (e.g., singletons)
&lt;/span&gt;&lt;span class="n"&gt;storage_client&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;storage&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nc"&gt;Client&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;
&lt;span class="n"&gt;bq_client&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;bigquery&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nc"&gt;Client&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;
&lt;span class="n"&gt;gcs_bucket&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;storage_client&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;bucket&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;GCS_BUCKET_NAME&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

&lt;span class="nd"&gt;@contextmanager&lt;/span&gt;
&lt;span class="k"&gt;def&lt;/span&gt; &lt;span class="nf"&gt;get_server_side_cursor&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;table_name&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;str&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;columns&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;list&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="nb"&gt;str&lt;/span&gt;&lt;span class="p"&gt;]):&lt;/span&gt;
    &lt;span class="sh"&gt;"""&lt;/span&gt;&lt;span class="s"&gt;
    Creates a server-side cursor to stream data from Postgres without
    loading it all into the client&lt;/span&gt;&lt;span class="sh"&gt;'&lt;/span&gt;&lt;span class="s"&gt;s memory. This is CRITICAL.
    &lt;/span&gt;&lt;span class="sh"&gt;"""&lt;/span&gt;
    &lt;span class="n"&gt;conn&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;psycopg&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;connect&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;PG_CONN_STRING&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;autocommit&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="bp"&gt;True&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="c1"&gt;# A server-side cursor is named. The client only ever holds a small buffer.
&lt;/span&gt;    &lt;span class="n"&gt;cursor_name&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;migration_cursor_&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;uuid&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;uuid4&lt;/span&gt;&lt;span class="p"&gt;().&lt;/span&gt;&lt;span class="nb"&gt;hex&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
    &lt;span class="n"&gt;cursor&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;conn&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;cursor&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;name&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;cursor_name&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="n"&gt;cursor&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;itersize&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;CHUNK_SIZE&lt;/span&gt;  &lt;span class="c1"&gt;# How many rows psycopg fetches at once internally
&lt;/span&gt;
    &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;INFO: Starting stream from table &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;table_name&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;...&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="k"&gt;try&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="n"&gt;column_str&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;, &lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;join&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;'"&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;c&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"'&lt;/span&gt; &lt;span class="k"&gt;for&lt;/span&gt; &lt;span class="n"&gt;c&lt;/span&gt; &lt;span class="ow"&gt;in&lt;/span&gt; &lt;span class="n"&gt;columns&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
        &lt;span class="n"&gt;cursor&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;execute&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;SELECT &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;column_str&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt; FROM &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;table_name&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
        &lt;span class="k"&gt;yield&lt;/span&gt; &lt;span class="n"&gt;cursor&lt;/span&gt;
    &lt;span class="k"&gt;finally&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;INFO: Closing connection.&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
        &lt;span class="n"&gt;cursor&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;close&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;
        &lt;span class="n"&gt;conn&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;close&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;

&lt;span class="k"&gt;def&lt;/span&gt; &lt;span class="nf"&gt;upload_to_gcs&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;file_path&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;str&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;gcs_path&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;str&lt;/span&gt;&lt;span class="p"&gt;):&lt;/span&gt;
    &lt;span class="sh"&gt;"""&lt;/span&gt;&lt;span class="s"&gt;Uploads a local file to GCS.&lt;/span&gt;&lt;span class="sh"&gt;"""&lt;/span&gt;
    &lt;span class="n"&gt;blob&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;gcs_bucket&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;blob&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;gcs_path&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="n"&gt;blob&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;upload_from_filename&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;file_path&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;INFO: Uploaded &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;file_path&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt; to gs://&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;GCS_BUCKET_NAME&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;/&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;gcs_path&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

&lt;span class="k"&gt;def&lt;/span&gt; &lt;span class="nf"&gt;stream_pg_to_gcs_parquet&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;table_name&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;str&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;columns&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;list&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="nb"&gt;str&lt;/span&gt;&lt;span class="p"&gt;]):&lt;/span&gt;
    &lt;span class="sh"&gt;"""&lt;/span&gt;&lt;span class="s"&gt;
    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.
    &lt;/span&gt;&lt;span class="sh"&gt;"""&lt;/span&gt;
    &lt;span class="n"&gt;gcs_paths&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="p"&gt;[]&lt;/span&gt;
    &lt;span class="k"&gt;with&lt;/span&gt; &lt;span class="nf"&gt;get_server_side_cursor&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;table_name&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;columns&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="k"&gt;as&lt;/span&gt; &lt;span class="n"&gt;cursor&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="n"&gt;chunk_num&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="mi"&gt;0&lt;/span&gt;
        &lt;span class="k"&gt;while&lt;/span&gt; &lt;span class="bp"&gt;True&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
            &lt;span class="c1"&gt;# fetchmany is the key to streaming.
&lt;/span&gt;            &lt;span class="n"&gt;rows&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;cursor&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;fetchmany&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;CHUNK_SIZE&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
            &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="ow"&gt;not&lt;/span&gt; &lt;span class="n"&gt;rows&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
                &lt;span class="k"&gt;break&lt;/span&gt;

            &lt;span class="n"&gt;chunk_num&lt;/span&gt; &lt;span class="o"&gt;+=&lt;/span&gt; &lt;span class="mi"&gt;1&lt;/span&gt;
            &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;INFO: Processing chunk &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;chunk_num&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt; with &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="nf"&gt;len&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;rows&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt; rows...&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

            &lt;span class="c1"&gt;# Zero-copy conversion to an Arrow Table
&lt;/span&gt;            &lt;span class="c1"&gt;# PyArrow can often directly map Postgres memory to Arrow memory.
&lt;/span&gt;            &lt;span class="k"&gt;try&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
                &lt;span class="n"&gt;record_batch&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;pa&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;RecordBatch&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;from_pylist&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;
                    &lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="nf"&gt;dict&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="nf"&gt;zip&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;columns&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;row&lt;/span&gt;&lt;span class="p"&gt;))&lt;/span&gt; &lt;span class="k"&gt;for&lt;/span&gt; &lt;span class="n"&gt;row&lt;/span&gt; &lt;span class="ow"&gt;in&lt;/span&gt; &lt;span class="n"&gt;rows&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;
                &lt;span class="p"&gt;)&lt;/span&gt;
                &lt;span class="n"&gt;arrow_table&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;pa&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;Table&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;from_batches&lt;/span&gt;&lt;span class="p"&gt;([&lt;/span&gt;&lt;span class="n"&gt;record_batch&lt;/span&gt;&lt;span class="p"&gt;])&lt;/span&gt;
            &lt;span class="k"&gt;except&lt;/span&gt; &lt;span class="nb"&gt;Exception&lt;/span&gt; &lt;span class="k"&gt;as&lt;/span&gt; &lt;span class="n"&gt;e&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
                &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;ERROR: Could not convert chunk to Arrow: &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;e&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
                &lt;span class="c1"&gt;# Add dead-letter queue logic here
&lt;/span&gt;                &lt;span class="k"&gt;continue&lt;/span&gt;

            &lt;span class="c1"&gt;# Write to a local Parquet file
&lt;/span&gt;            &lt;span class="n"&gt;local_filename&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;/tmp/&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;table_name&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;_chunk_&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;chunk_num&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;.parquet&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
            &lt;span class="n"&gt;pq&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;write_table&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;arrow_table&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;local_filename&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;compression&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="sh"&gt;'&lt;/span&gt;&lt;span class="s"&gt;SNAPPY&lt;/span&gt;&lt;span class="sh"&gt;'&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

            &lt;span class="c1"&gt;# Upload to GCS
&lt;/span&gt;            &lt;span class="n"&gt;gcs_path&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;GCS_PREFIX&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;/&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;table_name&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;/part-&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;chunk_num&lt;/span&gt;&lt;span class="si"&gt;:&lt;/span&gt;&lt;span class="mi"&gt;05&lt;/span&gt;&lt;span class="n"&gt;d&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;.parquet&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
            &lt;span class="nf"&gt;upload_to_gcs&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;local_filename&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;gcs_path&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
            &lt;span class="n"&gt;gcs_paths&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;append&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;gs://&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;GCS_BUCKET_NAME&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;/&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;gcs_path&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

            &lt;span class="n"&gt;os&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;remove&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;local_filename&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

    &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;INFO: Stream finished.&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;gcs_paths&lt;/span&gt;

&lt;span class="k"&gt;def&lt;/span&gt; &lt;span class="nf"&gt;load_gcs_to_bigquery&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;gcs_paths&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;list&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="nb"&gt;str&lt;/span&gt;&lt;span class="p"&gt;]):&lt;/span&gt;
    &lt;span class="sh"&gt;"""&lt;/span&gt;&lt;span class="s"&gt;Kicks off a BigQuery load job from GCS.&lt;/span&gt;&lt;span class="sh"&gt;"""&lt;/span&gt;
    &lt;span class="n"&gt;job_config&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;bigquery&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nc"&gt;LoadJobConfig&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;
        &lt;span class="n"&gt;source_format&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;bigquery&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;SourceFormat&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;PARQUET&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
        &lt;span class="n"&gt;write_disposition&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;bigquery&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;WriteDisposition&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;WRITE_TRUNCATE&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="c1"&gt;# Overwrite table
&lt;/span&gt;    &lt;span class="p"&gt;)&lt;/span&gt;

    &lt;span class="n"&gt;table_ref&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;bq_client&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;project&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;.&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;BQ_DATASET_ID&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;.&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;BQ_TABLE_ID&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
    &lt;span class="n"&gt;load_job&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;bq_client&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;load_table_from_uri&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;gcs_paths&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;table_ref&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;job_config&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;job_config&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;INFO: Starting BigQuery load job &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;load_job&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;job_id&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

    &lt;span class="n"&gt;load_job&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;result&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt; &lt;span class="c1"&gt;# Waits for the job to complete.
&lt;/span&gt;    &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;INFO: BigQuery load job finished.&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

&lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="n"&gt;__name__&lt;/span&gt; &lt;span class="o"&gt;==&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;__main__&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
    &lt;span class="c1"&gt;# In a real scenario, you'd get these columns from information_schema
&lt;/span&gt;    &lt;span class="c1"&gt;# and handle data type mapping carefully!
&lt;/span&gt;    &lt;span class="n"&gt;table_columns&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;id&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;user_id&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;event_type&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;event_timestamp&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;payload&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;

    &lt;span class="n"&gt;gcs_file_paths&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nf"&gt;stream_pg_to_gcs_parquet&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;events&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;table_columns&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

    &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="n"&gt;gcs_file_paths&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="nf"&gt;load_gcs_to_bigquery&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;gcs_file_paths&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h3&gt;
  
  
  Phase 2 &amp;amp; 3: The CDC Stream &amp;amp; The Terrifying Cutover
&lt;/h3&gt;

&lt;p&gt;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 &lt;code&gt;MERGE&lt;/code&gt; statements.&lt;/p&gt;

&lt;p&gt;This architecture diagram shows the two paths running in parallel:&lt;/p&gt;

&lt;p&gt;
  &lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fraw.githubusercontent.com%2Fjun-matsui%2Foutscape%2Fmain%2Fdocs%2Fassets%2Fdiagrams%2Fpipeline_sequence.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fraw.githubusercontent.com%2Fjun-matsui%2Foutscape%2Fmain%2Fdocs%2Fassets%2Fdiagrams%2Fpipeline_sequence.png" alt="Database Migration Architecture and Pipeline Sequence Diagram" width="800" height="494"&gt;&lt;/a&gt;
&lt;/p&gt;

&lt;p&gt;
  &lt;a href="https://mermaid.live/view#pako:bZRNb9pAEIb_ysinRAIlbdWLD5FqIEgVTU0dKVKVy2AP9gp71zFqFPW_981a2kZqTz1s-Hk-zwwzs2Bih2IGxL5w5-1EkMc7mzbAy-rHKk4b6zI7Z2vgw3rDNuE_q4qZvZ0lE5vVq7V-eUoW1XgH8529gXkLp3u3Pq0WfO8a513f5M3iP2Yp-3k2r7q6gV6322a36-r9o-9Vb_v-j0P39X6v2vX5vj0k9z_Nrmk-k6q97s18b5_u-9_h7Hvb2-1zXy9O1e2_8Fh_iV1VAM8v5Q7evv27_11-7v5U-7-b4v7X67e7_X5_1V8WwPPy8fHxfH1f21-7_f4f" rel="noopener noreferrer"&gt;🔍 &lt;b&gt;Click to View &amp;amp; Pan/Zoom in Mermaid Live Interactive Editor ↗&lt;/b&gt;&lt;/a&gt;
&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;The Moment of Truth: Data Integrity Validation&lt;/strong&gt;&lt;/p&gt;

&lt;p&gt;You &lt;em&gt;cannot&lt;/em&gt; skip this step. Hope is not a strategy. Before flipping the DNS, you must mathematically prove that the data is identical.&lt;/p&gt;

&lt;p&gt;🔥 &lt;strong&gt;Seniority Check:&lt;/strong&gt; Do you know why &lt;code&gt;COUNT(*)&lt;/code&gt; is not enough? Because &lt;code&gt;UPDATE&lt;/code&gt; operations can change data without changing the row count.&lt;/p&gt;

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

&lt;p&gt;Here's a simple Python validator:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="c1"&gt;# A simplified validation script
&lt;/span&gt;&lt;span class="k"&gt;def&lt;/span&gt; &lt;span class="nf"&gt;validate_data_integrity&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;pg_table&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;bq_table&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;column&lt;/span&gt;&lt;span class="p"&gt;):&lt;/span&gt;
    &lt;span class="sh"&gt;"""&lt;/span&gt;&lt;span class="s"&gt;Compares checksums between Postgres and BigQuery.&lt;/span&gt;&lt;span class="sh"&gt;"""&lt;/span&gt;

    &lt;span class="c1"&gt;# --- Postgres Checksum ---
&lt;/span&gt;    &lt;span class="k"&gt;with&lt;/span&gt; &lt;span class="n"&gt;psycopg&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;connect&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;PG_CONN_STRING&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="k"&gt;as&lt;/span&gt; &lt;span class="n"&gt;conn&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="k"&gt;with&lt;/span&gt; &lt;span class="n"&gt;conn&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;cursor&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt; &lt;span class="k"&gt;as&lt;/span&gt; &lt;span class="n"&gt;cur&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
            &lt;span class="c1"&gt;# BIT_XOR is a great way to 'fingerprint' a numeric column
&lt;/span&gt;            &lt;span class="n"&gt;cur&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;execute&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;SELECT BIT_XOR(CAST(&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;column&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt; AS BIGINT)) FROM &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;pg_table&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;;&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
            &lt;span class="n"&gt;pg_checksum&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;cur&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;fetchone&lt;/span&gt;&lt;span class="p"&gt;()[&lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;
            &lt;span class="n"&gt;cur&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;execute&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;SELECT COUNT(*) FROM &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;pg_table&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;;&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
            &lt;span class="n"&gt;pg_count&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;cur&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;fetchone&lt;/span&gt;&lt;span class="p"&gt;()[&lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;

    &lt;span class="c1"&gt;# --- BigQuery Checksum ---
&lt;/span&gt;    &lt;span class="n"&gt;bq_client&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;bigquery&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nc"&gt;Client&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;
    &lt;span class="c1"&gt;# BigQuery's equivalent for BIT_XOR
&lt;/span&gt;    &lt;span class="n"&gt;query&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"""&lt;/span&gt;&lt;span class="s"&gt;
        SELECT BIT_XOR(CAST(&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;column&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt; AS INT64)) as checksum, COUNT(*) as count
        FROM `&lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;bq_table&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;`
    &lt;/span&gt;&lt;span class="sh"&gt;"""&lt;/span&gt;
    &lt;span class="n"&gt;query_job&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;bq_client&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;query&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;query&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="n"&gt;results&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;query_job&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;result&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;
    &lt;span class="k"&gt;for&lt;/span&gt; &lt;span class="n"&gt;row&lt;/span&gt; &lt;span class="ow"&gt;in&lt;/span&gt; &lt;span class="n"&gt;results&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="n"&gt;bq_checksum&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;row&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;checksum&lt;/span&gt;
        &lt;span class="n"&gt;bq_count&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;row&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;count&lt;/span&gt;

    &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;--- DATA INTEGRITY CHECK ---&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;Postgres Count: &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;pg_count&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;, BigQuery Count: &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;bq_count&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;f&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;Postgres XOR:   &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;pg_checksum&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="s"&gt;, BigQuery XOR:   &lt;/span&gt;&lt;span class="si"&gt;{&lt;/span&gt;&lt;span class="n"&gt;bq_checksum&lt;/span&gt;&lt;span class="si"&gt;}&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;

    &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="n"&gt;pg_count&lt;/span&gt; &lt;span class="o"&gt;==&lt;/span&gt; &lt;span class="n"&gt;bq_count&lt;/span&gt; &lt;span class="ow"&gt;and&lt;/span&gt; &lt;span class="n"&gt;pg_checksum&lt;/span&gt; &lt;span class="o"&gt;==&lt;/span&gt; &lt;span class="n"&gt;bq_checksum&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;✅ SUCCESS: Data is consistent!&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
        &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="bp"&gt;True&lt;/span&gt;
    &lt;span class="k"&gt;else&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="nf"&gt;print&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;❌ FAILURE: Data mismatch detected!&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
        &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="bp"&gt;False&lt;/span&gt;

&lt;span class="c1"&gt;# Run this right before the final cutover
# validate_data_integrity("public.events", "your_project.your_dataset.your_migrated_table", "id")
&lt;/span&gt;&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h3&gt;
  
  
  The "Gotchas" The Docs Don't Tell You
&lt;/h3&gt;

&lt;p&gt;This was not a smooth ride. Here are the landmines we hit:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;&lt;p&gt;&lt;strong&gt;Timezones are the Devil's Work:&lt;/strong&gt; Postgres &lt;code&gt;TIMESTAMPTZ&lt;/code&gt; stores timestamps in UTC internally. BigQuery's &lt;code&gt;TIMESTAMP&lt;/code&gt; 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 &lt;em&gt;must&lt;/em&gt; audit every single timestamp column and ensure your ETL pipeline explicitly sets the timezone to UTC.&lt;/p&gt;&lt;/li&gt;
&lt;li&gt;&lt;p&gt;&lt;strong&gt;The Hidden Costs of GCS:&lt;/strong&gt; 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.&lt;/p&gt;&lt;/li&gt;
&lt;li&gt;&lt;p&gt;&lt;strong&gt;BigQuery &lt;code&gt;MERGE&lt;/code&gt; Limits:&lt;/strong&gt; BigQuery's &lt;code&gt;MERGE&lt;/code&gt; 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 &lt;code&gt;MERGE&lt;/code&gt; operation every few minutes, not on every single event.&lt;/p&gt;&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;🕹️ &lt;b&gt;Mini-Quiz: A query that was fast on Postgres is suddenly slow on BigQuery. What's the most likely culprit?&lt;/b&gt; (Click to reveal)&lt;/p&gt;

&lt;p&gt;&amp;gt; &lt;strong&gt;The Culprit:&lt;/strong&gt; You're likely doing a &lt;code&gt;SELECT *&lt;/code&gt; 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 &lt;code&gt;SELECT *&lt;/code&gt; forces a full scan of &lt;em&gt;every single column&lt;/em&gt;, even if you only need two of them. It completely negates the columnar advantage.&lt;br&gt;
&amp;gt;&lt;br&gt;
&amp;gt; &lt;strong&gt;The Fix:&lt;/strong&gt; Be explicit with &lt;code&gt;SELECT column1, column2&lt;/code&gt;. And critically, partition your BigQuery table (e.g., by &lt;code&gt;event_timestamp&lt;/code&gt;) and cluster it by commonly filtered columns (e.g., &lt;code&gt;user_id&lt;/code&gt;). 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.&lt;/p&gt;

&lt;h3&gt;
  
  
  The Payoff: From Pager Alarms to Peaceful Sleep
&lt;/h3&gt;

&lt;p&gt;After the cutover, the difference was night and day.&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Metric&lt;/th&gt;
&lt;th&gt;Before (PostgreSQL)&lt;/th&gt;
&lt;th&gt;After (BigQuery)&lt;/th&gt;
&lt;th&gt;Impact&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;P95 Query Latency&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;45 minutes (or timeout)&lt;/td&gt;
&lt;td&gt;2.1 seconds&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;-99.99%&lt;/strong&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;P50 Query Latency&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;38 seconds&lt;/td&gt;
&lt;td&gt;650 milliseconds&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;-98.2%&lt;/strong&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;Concurrent Analysts&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;~3 (before locking)&lt;/td&gt;
&lt;td&gt;100+&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;+33x&lt;/strong&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;DBA Heart Rate&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;140 bpm (on-call)&lt;/td&gt;
&lt;td&gt;65 bpm (sleeping)&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;-53.5%&lt;/strong&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;Storage Cost&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;~$2,500/mo (Managed PG)&lt;/td&gt;
&lt;td&gt;~$300/mo (BQ Storage)&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;-88%&lt;/strong&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;The PagerDuty alerts for &lt;code&gt;iowait&lt;/code&gt; 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.&lt;/p&gt;

&lt;p&gt;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, &lt;em&gt;always&lt;/em&gt; validate your data.&lt;/p&gt;




&lt;h3&gt;
  
  
  Interactive Troubleshooting Guide
&lt;/h3&gt;

&lt;p&gt;Facing your own migration woes? Use this decision tree to diagnose the bottleneck.&lt;/p&gt;

&lt;p&gt;
  &lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fraw.githubusercontent.com%2Fjun-matsui%2Foutscape%2Fmain%2Fdocs%2Fassets%2Fdiagrams%2Fdecision_tree.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fraw.githubusercontent.com%2Fjun-matsui%2Foutscape%2Fmain%2Fdocs%2Fassets%2Fdiagrams%2Fdecision_tree.png" alt="Migration Troubleshooting Decision Tree Diagram" width="800" height="1021"&gt;&lt;/a&gt;
&lt;/p&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;┌────────────────────────────────────────────────────────┐
│        🚨 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!
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;
  &lt;a href="https://mermaid.live/view#pako:dZNdb5swFIb_ypGvh6re9mJVg5m6LV9TGk3TqCrXnIBbwOzYLIvS_PcdAyUs0rgBmZeH9zmYo9A2Q3EDYlfavS4UeXiQaQ183P1MxcLkpLyxNWz4PliCT8qUps5vU_EIUfQRZsdUxAXqV9jYljSCnKXi1BNmfeKaQfcmLyBeb68-r4Bpa-t8TugCZsheh_DbD3RvfP3EZ35qo1WJ0DZA2JRGq1DAF2S95-XkYQ6uQcwuGUvbIZY2EPpSxsHO1NglB7uuWzy2D7Tvll6RxvpxHxnrL7CydICtUznCFaxWC_hqyhKzs0U8sYjfLbYOwSH9RoqcyRB0S86SA1VnsD7cEfFknSdUFQ_2EhVk4kGm7_ePTEhOheQotES_5_xoI_v7wWbOLwxjbPOiaT08cxKxBl3aNuta2TpqCKuzl5x4yanXplZNc4i0rfgB55CNFP1qkanK62IqJM9CchAaOv738ySjzczk31rk6c-tyuCLfR69kj4ZvHgZdv3-hL3xBThdYKUAecQ02WvJxCZ5t5HKK_CHBqEyrgrlbyDDUAvwT9h9xg-8S04wSgajrt4L95gqiQ8gKqRKmYx_taPwTOl-OuartvTidPoL" rel="noopener noreferrer"&gt;🔍 &lt;b&gt;Click to View &amp;amp; Pan/Zoom in Mermaid Live Interactive Editor ↗&lt;/b&gt;&lt;/a&gt;
&lt;/p&gt;

&lt;h3&gt;
  
  
  Further Reading &amp;amp; Official Docs
&lt;/h3&gt;

&lt;ul&gt;
&lt;li&gt;  &lt;strong&gt;Python:&lt;/strong&gt;

&lt;ul&gt;
&lt;li&gt;  &lt;a href="https://www.psycopg.org/psycopg3/docs/advanced/cursors.html#server-side-cursors" rel="noopener noreferrer"&gt;Psycopg 3 Documentation (Server-side cursors)&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;  &lt;a href="https://arrow.apache.org/docs/python/generated/pyarrow.RecordBatch.html" rel="noopener noreferrer"&gt;PyArrow &lt;code&gt;RecordBatch&lt;/code&gt; Documentation&lt;/a&gt;
&lt;/li&gt;
&lt;/ul&gt;
&lt;/li&gt;
&lt;li&gt;  &lt;strong&gt;Google Cloud:&lt;/strong&gt;

&lt;ul&gt;
&lt;li&gt;  &lt;a href="https://cloud.google.com/bigquery/docs/loading-data-cloud-storage-parquet" rel="noopener noreferrer"&gt;BigQuery: Loading data from Parquet&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;  &lt;a href="https://cloud.google.com/bigquery/docs/reference/standard-sql/dml-syntax#merge_statement" rel="noopener noreferrer"&gt;BigQuery: &lt;code&gt;MERGE&lt;/code&gt; statement DML&lt;/a&gt;
&lt;/li&gt;
&lt;/ul&gt;
&lt;/li&gt;
&lt;li&gt;  &lt;strong&gt;PostgreSQL:&lt;/strong&gt;

&lt;ul&gt;
&lt;li&gt;  &lt;a href="https://www.postgresql.org/docs/current/logical-replication.html" rel="noopener noreferrer"&gt;Logical Replication Documentation&lt;/a&gt;
&lt;/li&gt;
&lt;/ul&gt;
&lt;/li&gt;
&lt;/ul&gt;

</description>
      <category>python</category>
      <category>database</category>
      <category>cloud</category>
      <category>architecture</category>
    </item>
    <item>
      <title>Benchmarking LLMs on ETL Logic Synthesis: Can AI Truly Replace Data Pipeline Scripting?</title>
      <dc:creator>Jun Matsui</dc:creator>
      <pubDate>Thu, 08 Oct 2026 21:06:42 +0000</pubDate>
      <link>https://dev.to/jun-matsui/benchmarking-llms-on-etl-logic-synthesis-can-ai-truly-replace-data-pipeline-scripting-27ln</link>
      <guid>https://dev.to/jun-matsui/benchmarking-llms-on-etl-logic-synthesis-can-ai-truly-replace-data-pipeline-scripting-27ln</guid>
      <description>&lt;p&gt;&lt;em&gt;This is a submission for the &lt;a href="https://dev.to/challenges/kaggle-2026-09-23"&gt;Kaggle Benchmarking Challenge&lt;/a&gt;&lt;/em&gt;&lt;/p&gt;




&lt;h2&gt;
  
  
  What I Benchmarked
&lt;/h2&gt;

&lt;p&gt;In enterprise data engineering, migrating visual ETL pipelines (from tools like &lt;strong&gt;Knime&lt;/strong&gt;, &lt;strong&gt;Alteryx&lt;/strong&gt;, or &lt;strong&gt;SSIS&lt;/strong&gt;) or complex business pseudocode into performant, vectorized Python (&lt;code&gt;pandas&lt;/code&gt; / &lt;code&gt;polars&lt;/code&gt;) is one of the most critical and recurring challenges.&lt;/p&gt;

&lt;p&gt;While standard benchmarks evaluate generic programming puzzles or synthetic LeetCode algorithms, real-world data pipelines break due to subtle edge cases. I built the &lt;strong&gt;ETL-to-Python Code Synthesis Benchmark&lt;/strong&gt; to evaluate whether LLMs can synthesize clean, idiomatic, and robust Python code from visual workflow specifications.&lt;br&gt;
&lt;/p&gt;

&lt;pre data-lang="mermaid"&gt;&lt;code&gt;flowchart TD
    Start["🚨 Input: Legacy Visual ETL Node Graph"] --&amp;gt; T1["Task 01: Left Join &amp;amp; Imputation&amp;lt;br/&amp;gt;• Coerce nulls&amp;lt;br/&amp;gt;• Calculate is_vip flag"]
    Start --&amp;gt; T2["Task 02: Regex Extraction&amp;lt;br/&amp;gt;• Parse key-value logs&amp;lt;br/&amp;gt;• Retain corrupted rows"]
    Start --&amp;gt; T3["Task 03: Cumulative Windows&amp;lt;br/&amp;gt;• Running total cumsum()&amp;lt;br/&amp;gt;• Intra-department rank"]
    Start --&amp;gt; T4["Task 04: Matrix Reshaping&amp;lt;br/&amp;gt;• Melt wide quarters&amp;lt;br/&amp;gt;• Flatten MultiIndex headers"]

    T1 --&amp;gt; Sandbox["🧪 Sandboxed PyTest Execution Engine"]
    T2 --&amp;gt; Sandbox
    T3 --&amp;gt; Sandbox
    T4 --&amp;gt; Sandbox
    Sandbox --&amp;gt; Leaderboard["🏆 Sub-millisecond DataFrame Assertion Leaderboard"]

    style Start fill:#1e1e2e,stroke:#89b4fa,color:#cdd6f4
    style Sandbox fill:#313244,stroke:#f9e2af,color:#cdd6f4
    style Leaderboard fill:#14532d,stroke:#22c55e,color:#f0fdf4&lt;/code&gt;&lt;/pre&gt;



&lt;h3&gt;
  
  
  The 4 Evaluated Tasks:
&lt;/h3&gt;

&lt;ol&gt;
&lt;li&gt;
&lt;strong&gt;&lt;code&gt;etl_01&lt;/code&gt; (Joiner &amp;amp; Missing Value Imputation with Type Coercion):&lt;/strong&gt; Relational left joins with unmapped keys, safe type casting for corrupt numeric values, and conditional multi-column business flags (&lt;code&gt;is_vip&lt;/code&gt;).&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;&lt;code&gt;etl_02&lt;/code&gt; (Regex Extractor &amp;amp; Multi-Column Sanitizer):&lt;/strong&gt; Parsing semi-structured key-value log entries, handling malformed/corrupted rows without throwing exceptions, and applying strict exclusionary filtering.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;&lt;code&gt;etl_03&lt;/code&gt; (GroupLoop to Vectorized Cumulative Windows):&lt;/strong&gt; Eliminating slow iterative loops by synthesizing vectorized cumulative sums (&lt;code&gt;.cumsum()&lt;/code&gt;), target achievement ratios, rolling 3-month averages, and intra-department dense rankings.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;&lt;code&gt;etl_04&lt;/code&gt; (Unpivoting, Pivoting &amp;amp; Multi-Level Column Flattening):&lt;/strong&gt; Reshaping wide multi-quarter tables via melting, splitting composite temporal strings, pivoting, and flattening complex MultiIndex column headers to single-level snake_case schemas.&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;Each task runs inside an automated Python sandbox that tests DataFrame structural integrity, exact type fidelity, and output values under sub-millisecond execution times.&lt;/p&gt;




&lt;h2&gt;
  
  
  Models Tested
&lt;/h2&gt;

&lt;p&gt;I evaluated modern state-of-the-art models from Google DeepMind under deterministic zero-shot settings (&lt;code&gt;temperature = 0.0&lt;/code&gt;):&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;&lt;code&gt;gemini-3.8-flash&lt;/code&gt;&lt;/strong&gt;: Picked to evaluate high-throughput, low-latency code synthesis for real-time developer tooling and data transpilers.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;&lt;code&gt;gemini-2.5-pro&lt;/code&gt;&lt;/strong&gt;: Picked to evaluate deep multi-step reasoning capabilities when faced with intricate analytical requirements.&lt;/li&gt;
&lt;/ul&gt;




&lt;h2&gt;
  
  
  Findings
&lt;/h2&gt;

&lt;h3&gt;
  
  
  📊 Benchmark Leaderboard
&lt;/h3&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Model&lt;/th&gt;
&lt;th&gt;Accuracy (Passed / Total)&lt;/th&gt;
&lt;th&gt;Avg Score&lt;/th&gt;
&lt;th&gt;Avg API Latency&lt;/th&gt;
&lt;th&gt;Sandbox Assertion Speed&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;🥇 &lt;strong&gt;&lt;code&gt;gemini-3.8-flash&lt;/code&gt;&lt;/strong&gt;
&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;100.0% (4/4)&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;1.00&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;14.89 s&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;~13.4 ms&lt;/strong&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;🥈 &lt;strong&gt;&lt;code&gt;gemini-2.5-pro&lt;/code&gt;&lt;/strong&gt;
&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;100.0% (4/4)&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;1.00&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;33.78 s&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;&lt;strong&gt;~16.0 ms&lt;/strong&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;h3&gt;
  
  
  Breakdown by Task:
&lt;/h3&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Task ID&lt;/th&gt;
&lt;th&gt;Description&lt;/th&gt;
&lt;th&gt;&lt;code&gt;gemini-3.8-flash&lt;/code&gt;&lt;/th&gt;
&lt;th&gt;&lt;code&gt;gemini-2.5-pro&lt;/code&gt;&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;&lt;code&gt;etl_01&lt;/code&gt;&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;Left Join, Missing Values &amp;amp; Type Coercion&lt;/td&gt;
&lt;td&gt;✅ &lt;strong&gt;PASS&lt;/strong&gt;
&lt;/td&gt;
&lt;td&gt;✅ &lt;strong&gt;PASS&lt;/strong&gt;
&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;&lt;code&gt;etl_02&lt;/code&gt;&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;Regex Extraction &amp;amp; Edge-Case Sanitization&lt;/td&gt;
&lt;td&gt;✅ &lt;strong&gt;PASS&lt;/strong&gt;
&lt;/td&gt;
&lt;td&gt;✅ &lt;strong&gt;PASS&lt;/strong&gt;
&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;&lt;code&gt;etl_03&lt;/code&gt;&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;Cumulative Windows &amp;amp; Department Ranks&lt;/td&gt;
&lt;td&gt;✅ &lt;strong&gt;PASS&lt;/strong&gt;
&lt;/td&gt;
&lt;td&gt;✅ &lt;strong&gt;PASS&lt;/strong&gt;
&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;&lt;code&gt;etl_04&lt;/code&gt;&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;Matrix Reshaping &amp;amp; MultiIndex Flattening&lt;/td&gt;
&lt;td&gt;✅ &lt;strong&gt;PASS&lt;/strong&gt;
&lt;/td&gt;
&lt;td&gt;✅ &lt;strong&gt;PASS&lt;/strong&gt;
&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;h3&gt;
  
  
  🔍 Main Insights &amp;amp; Surprises
&lt;/h3&gt;

&lt;ol&gt;
&lt;li&gt;
&lt;strong&gt;The "Corrupted Row" Trap in Log Parsing (&lt;code&gt;etl_02&lt;/code&gt;):&lt;/strong&gt;
When parsing log streams that mix structured records with unformatted, corrupted strings (e.g. &lt;code&gt;"CORRUPTED_LINE_WITHOUT_DELIMITERS"&lt;/code&gt;), models often default to chaining aggressive &lt;code&gt;dropna()&lt;/code&gt; operations that delete the entire corrupted line.

&lt;ul&gt;
&lt;li&gt;
&lt;em&gt;The Key Insight:&lt;/em&gt; High-performing code synthesis separates &lt;em&gt;sanitization&lt;/em&gt; from &lt;em&gt;filtering&lt;/em&gt;, extracting named capture groups with &lt;code&gt;.fillna("anonymous")&lt;/code&gt; fallbacks to avoid silent audit data loss.&lt;/li&gt;
&lt;/ul&gt;
&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;🕹️ &lt;b&gt;Mini-Quiz: Why is df.iterrows() the enemy of production ETL pipelines?&lt;/b&gt; (Click to reveal)&lt;/p&gt;

&lt;p&gt;&amp;gt; &lt;strong&gt;The Cost:&lt;/strong&gt; Iterating over DataFrame rows with &lt;code&gt;for index, row in df.iterrows()&lt;/code&gt; converts each row into a pandas Series, creating massive Python overhead and slowing execution by up to &lt;strong&gt;100x–500x&lt;/strong&gt; compared to vectorized C-level operations like &lt;code&gt;df.groupby().cumsum()&lt;/code&gt; or &lt;code&gt;.rolling()&lt;/code&gt;.&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;&lt;p&gt;&lt;strong&gt;Native Loop Vectorization is Solved:&lt;/strong&gt;&lt;br&gt;
In Task 3 (translating Knime's iterative GroupLoop node), both models entirely avoided &lt;code&gt;for row in df.iterrows()&lt;/code&gt; or iterative Python loops. Both synthesized clean, vectorized &lt;code&gt;df.groupby('employee_id')['revenue'].cumsum()&lt;/code&gt; and &lt;code&gt;df.groupby('department')['revenue'].rank(ascending=False, method='min')&lt;/code&gt;, demonstrating strong intrinsic understanding of pandas performance optimization.&lt;/p&gt;&lt;/li&gt;
&lt;li&gt;&lt;p&gt;&lt;strong&gt;Flash Delivers 2.27x Higher Throughput:&lt;/strong&gt;&lt;br&gt;
&lt;code&gt;gemini-3.8-flash&lt;/code&gt; achieved a &lt;strong&gt;perfect 100% score in an average of 14.89 seconds per task&lt;/strong&gt;, compared to &lt;strong&gt;33.78 seconds for &lt;code&gt;gemini-2.5-pro&lt;/code&gt;&lt;/strong&gt;. For real-time IDE extensions and automated transpilers, Flash is clearly the most cost-effective choice.&lt;/p&gt;&lt;/li&gt;
&lt;/ol&gt;

&lt;h3&gt;
  
  
  🔮 What I Would Measure Next
&lt;/h3&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;Polars &amp;amp; PySpark Syntheses:&lt;/strong&gt; Evaluating whether LLMs can synthesize zero-copy LazyFrame queries in Polars with equal reliability.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;SQL Dialect Transpilation:&lt;/strong&gt; Measuring cross-engine translation from Snowflake SQL to Google Cloud BigQuery.&lt;/li&gt;
&lt;/ul&gt;




&lt;h2&gt;
  
  
  My Benchmark
&lt;/h2&gt;

&lt;p&gt;You can inspect, fork, and run this benchmark directly on Kaggle and GitHub:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;🔗 &lt;strong&gt;Kaggle Benchmark Notebook:&lt;/strong&gt; &lt;a href="https://www.kaggle.com/benchmarks" rel="noopener noreferrer"&gt;https://www.kaggle.com/benchmarks&lt;/a&gt; &lt;em&gt;(Task: &lt;code&gt;etl_knime_to_python_code_synthesis&lt;/code&gt;)&lt;/em&gt;
&lt;/li&gt;
&lt;li&gt;🐙 &lt;strong&gt;GitHub Repository:&lt;/strong&gt; &lt;a href="https://github.com/jun-matsui/kaggle-etl-benchmark" rel="noopener noreferrer"&gt;https://github.com/jun-matsui/kaggle-etl-benchmark&lt;/a&gt;
&lt;/li&gt;
&lt;/ul&gt;

&lt;h3&gt;
  
  
  Kaggle Task Implementation Snippet (&lt;code&gt;@kbench.task&lt;/code&gt;):
&lt;/h3&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;kbench&lt;/span&gt;
&lt;span class="kn"&gt;import&lt;/span&gt; &lt;span class="n"&gt;re&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;pandas&lt;/span&gt; &lt;span class="k"&gt;as&lt;/span&gt; &lt;span class="n"&gt;pd&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;numpy&lt;/span&gt; &lt;span class="k"&gt;as&lt;/span&gt; &lt;span class="n"&gt;np&lt;/span&gt;

&lt;span class="c1"&gt;# @kbench.task(
#     name="etl_knime_to_python_code_synthesis",
#     version="1.0.0",
#     description="Evaluates LLM capability in converting visual ETL pipeline logic into idiomatic, vectorized Python pandas code."
# )
&lt;/span&gt;&lt;span class="k"&gt;def&lt;/span&gt; &lt;span class="nf"&gt;evaluate_etl_benchmark&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;model_output&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;str&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;task_id&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;str&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;etl_01&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="nb"&gt;float&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
    &lt;span class="n"&gt;code_match&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;re&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;search&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sa"&gt;r&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;```

(?:python)?\s*(.*?)\s*

```&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;model_output&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;re&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;DOTALL&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="n"&gt;clean_code&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;code_match&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;group&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="p"&gt;).&lt;/span&gt;&lt;span class="nf"&gt;strip&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt; &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="n"&gt;code_match&lt;/span&gt; &lt;span class="k"&gt;else&lt;/span&gt; &lt;span class="n"&gt;model_output&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;strip&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;

    &lt;span class="n"&gt;local_scope&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;pd&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="n"&gt;pd&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;np&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="n"&gt;np&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;re&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="n"&gt;re&lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;
    &lt;span class="k"&gt;try&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="nf"&gt;exec&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;clean_code&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;local_scope&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;local_scope&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
        &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;transform_etl&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt; &lt;span class="ow"&gt;not&lt;/span&gt; &lt;span class="ow"&gt;in&lt;/span&gt; &lt;span class="n"&gt;local_scope&lt;/span&gt; &lt;span class="ow"&gt;or&lt;/span&gt; &lt;span class="ow"&gt;not&lt;/span&gt; &lt;span class="nf"&gt;callable&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;local_scope&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;transform_etl&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;]):&lt;/span&gt;
            &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="mf"&gt;0.0&lt;/span&gt;
        &lt;span class="c1"&gt;# Rigorous assertions on DataFrames
&lt;/span&gt;        &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="mf"&gt;1.0&lt;/span&gt;
    &lt;span class="k"&gt;except&lt;/span&gt; &lt;span class="nb"&gt;Exception&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;
        &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="mf"&gt;0.0&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;em&gt;All dataset fixtures, automated test suites, and runners are open-sourced at &lt;a href="https://github.com/jun-matsui/kaggle-etl-benchmark" rel="noopener noreferrer"&gt;github.com/jun-matsui/kaggle-etl-benchmark&lt;/a&gt;.&lt;/em&gt;&lt;/p&gt;

</description>
      <category>kagglechallenge</category>
      <category>ai</category>
      <category>python</category>
      <category>machinelearning</category>
    </item>
    <item>
      <title>Outscape – Disconnect from the Screen with Open-Weight AI</title>
      <dc:creator>Jun Matsui</dc:creator>
      <pubDate>Thu, 08 Oct 2026 18:30:49 +0000</pubDate>
      <link>https://dev.to/jun-matsui/outscape-disconnect-from-the-screen-with-open-weight-ai-4i27</link>
      <guid>https://dev.to/jun-matsui/outscape-disconnect-from-the-screen-with-open-weight-ai-4i27</guid>
      <description>&lt;p&gt;&lt;em&gt;This is a submission for the &lt;a href="https://dev.to/challenges/hacktoberfest-week1-2026-10-05"&gt;Hacktoberfest Open-Source AI Challenge Week 1: Touch Grass&lt;/a&gt;&lt;/em&gt;&lt;/p&gt;

&lt;h2&gt;
  
  
  What I Built
&lt;/h2&gt;

&lt;p&gt;&lt;strong&gt;Outscape&lt;/strong&gt; is an open-source, hands-free outdoor exploration companion powered by open-weight AI (&lt;strong&gt;Google Gemma 2&lt;/strong&gt;) with &lt;strong&gt;WebGPU on-device acceleration&lt;/strong&gt;. &lt;/p&gt;

&lt;p&gt;The goal is simple: &lt;strong&gt;make screen time the shortest part of your day&lt;/strong&gt;. Instead of staring at a mobile screen while walking through a park or hiking a forest ridge, Outscape generates real-time audio observations and mindful cues ("Touch Grass" prompts) directly to your headphones. You can put your phone in your pocket ("Pocket Mode") and let smart ambient audio guide your discovery of local trees, seasonal foliage, and songbird habitats.&lt;/p&gt;

&lt;p&gt;It works seamlessly on both &lt;strong&gt;mobile devices&lt;/strong&gt; (with real outdoor GPS tracking) and &lt;strong&gt;desktop PCs&lt;/strong&gt; (with an interactive walking simulator).&lt;/p&gt;

&lt;h2&gt;
  
  
  Demo
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;Live Application:&lt;/strong&gt; &lt;a href="https://jun-matsui.github.io/outscape/" rel="noopener noreferrer"&gt;https://jun-matsui.github.io/outscape/&lt;/a&gt;
&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Features in action:&lt;/strong&gt;

&lt;ul&gt;
&lt;li&gt;⚡ &lt;strong&gt;WebGPU On-Device Acceleration:&lt;/strong&gt; Leverages your device's physical GPU (Apple Silicon, NVIDIA, AMD) for local open-weight inference with zero telemetry.&lt;/li&gt;
&lt;li&gt;🎧 &lt;strong&gt;Pocket Mode:&lt;/strong&gt; Screen-dimmed, tranquil audio walk with background step tracking.&lt;/li&gt;
&lt;li&gt;🗺️ &lt;strong&gt;Interactive Nature Trails:&lt;/strong&gt; Maple &amp;amp; Oak Ridge (Fall foliage), Emerald Greenbelt (Urban flora), and Whispering Brook (Aquatic biome).&lt;/li&gt;
&lt;li&gt;☀️ &lt;strong&gt;Sunlight High-Contrast Mode:&lt;/strong&gt; Accessible outdoor visibility under bright sun.&lt;/li&gt;
&lt;li&gt;🚶‍♂️ &lt;strong&gt;Desktop Step Simulator:&lt;/strong&gt; Test trails and triggers with 1x, 2x, and 4x speed controls.&lt;/li&gt;
&lt;/ul&gt;
&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  Code
&lt;/h2&gt;


&lt;div class="ltag-github-readme-tag"&gt;
  &lt;div class="readme-overview"&gt;
    &lt;h2&gt;
      &lt;img src="https://assets.dev.to/assets/github-logo-5a155e1f9a670af7944dd5e12375bc76ed542ea80224905ecaf878b9157cdefc.svg" alt="GitHub logo"&gt;
      &lt;a href="https://github.com/jun-matsui" rel="noopener noreferrer"&gt;
        jun-matsui
      &lt;/a&gt; / &lt;a href="https://github.com/jun-matsui/outscape" rel="noopener noreferrer"&gt;
        outscape
      &lt;/a&gt;
    &lt;/h2&gt;
    &lt;h3&gt;
      Outscape — An open-source, hands-free outdoor exploration companion powered by open-weight AI. Designed to get people off their screens and into nature. Built for the Dev.to Hacktoberfest "Touch Grass" Challenge. Outscape — オープンウェイトAIを活用したハンズフリーの野外探索コンパニオン。画面から離れて自然と触れ合う体験を届けるオープンソースプロジェクト（Dev.to Hacktoberfest「Touch Grass」チャレンジ参加作品）。
    &lt;/h3&gt;
  &lt;/div&gt;
  &lt;div class="ltag-github-body"&gt;
    
&lt;div id="readme" class="md"&gt;&lt;div class="markdown-heading"&gt;
&lt;h1 class="heading-element"&gt;🌿 Outscape&lt;/h1&gt;
&lt;/div&gt;
&lt;p&gt;&lt;a href="https://dev.to/challenges/hacktoberfest-week1-2026-10-05" rel="nofollow"&gt;&lt;img src="https://camo.githubusercontent.com/26bc2681af54fed34934760edc66fe1d6db4a463e9feb9f1a0cb419f3b25e0ce/68747470733a2f2f696d672e736869656c64732e696f2f62616467652f4861636b746f626572666573742d323032362d6f72616e67653f7374796c653d666f722d7468652d6261646765266c6f676f3d6861636b746f62657266657374" alt="Hacktoberfest 2026"&gt;&lt;/a&gt;
&lt;a href="https://dev.to/challenges/hacktoberfest-week1-2026-10-05" rel="nofollow"&gt;&lt;img src="https://camo.githubusercontent.com/1f50f5b25f3716fdb5866b3f6b1563c5f632a51bbb0a199aad82c0639d4ba6f0/68747470733a2f2f696d672e736869656c64732e696f2f62616467652f4445562532304368616c6c656e67652d5765656b25323031253341253230546f75636825323047726173732d626c756576696f6c65743f7374796c653d666f722d7468652d6261646765" alt="DEV Challenge"&gt;&lt;/a&gt;
&lt;a href="https://github.com/jun-matsui/outscape/LICENSE" rel="noopener noreferrer"&gt;&lt;img src="https://camo.githubusercontent.com/9218332452902d9e542a100d0af126fd3174a116456614d2cf093546a13783db/68747470733a2f2f696d672e736869656c64732e696f2f62616467652f4c6963656e73652d4d49542d677265656e2e7376673f7374796c653d666f722d7468652d6261646765" alt="License: MIT"&gt;&lt;/a&gt;
&lt;a href="https://ai.google.dev/gemma" rel="nofollow noopener noreferrer"&gt;&lt;img src="https://camo.githubusercontent.com/ebe15724f78e9c8f83b17f0c42400df254ca649124a77796af94619cfa20f8db/68747470733a2f2f696d672e736869656c64732e696f2f62616467652f41492d4f70656e2d2d5765696768742532304d6f64656c732d626c75653f7374796c653d666f722d7468652d6261646765" alt="Open Source AI"&gt;&lt;/a&gt;&lt;/p&gt;
&lt;blockquote&gt;
&lt;p&gt;&lt;strong&gt;Disconnect from the screen. Reconnect with the wild.&lt;/strong&gt;&lt;br&gt;
&lt;strong&gt;画面を閉じて、自然とつながる。&lt;/strong&gt;&lt;/p&gt;
&lt;/blockquote&gt;

&lt;div class="markdown-heading"&gt;
&lt;h2 class="heading-element"&gt;📖 Overview / 概要&lt;/h2&gt;
&lt;/div&gt;
&lt;div class="markdown-heading"&gt;
&lt;h3 class="heading-element"&gt;🇬🇧 English&lt;/h3&gt;
&lt;/div&gt;
&lt;p&gt;&lt;strong&gt;Outscape&lt;/strong&gt; is an open-source outdoor navigation and nature exploration companion driven by open-weight AI. Designed specifically for the &lt;strong&gt;DEV.to Hacktoberfest "Touch Grass" Challenge&lt;/strong&gt;, Outscape aims to make screen interaction the shortest part of your day.&lt;/p&gt;
&lt;p&gt;Instead of keeping your eyes glued to a display while walking, hiking, or running, Outscape turns local environmental insights, seasonal highlights (like fall foliage or spring blooming), and trail discoveries into an immersive, &lt;strong&gt;hands-free audio experience&lt;/strong&gt;. Put your phone in your pocket, listen to smart audio cues, and explore the outdoors with your senses wide open.&lt;/p&gt;
&lt;div class="markdown-heading"&gt;
&lt;h3 class="heading-element"&gt;🇯🇵 日本語&lt;/h3&gt;

&lt;/div&gt;
&lt;p&gt;&lt;strong&gt;Outscape&lt;/strong&gt; は、オープンウェイト（オープンソース）AIを活用したハンズフリーの野外探索・音声ガイドコンパニオンです。&lt;strong&gt;DEV.to Hacktoberfest「Touch Grass」チャレンジ&lt;/strong&gt;に向けて開発されました。&lt;/p&gt;
&lt;p&gt;散歩やハイキング、ランニング中に画面を見つめ続けるのではなく、最小限の画面操作とスマートな音声ガイダンスによって、周囲の自然やトレイルの魅力を五感で体験できるよう設計されています。スマートフォンをポケットにしまい、自然の音に耳を傾けながら、外の世界へ踏み出しましょう。&lt;/p&gt;

&lt;div class="markdown-heading"&gt;
&lt;h2 class="heading-element"&gt;✨ Key Features / 主な機能&lt;/h2&gt;

&lt;/div&gt;
&lt;ul&gt;
&lt;li&gt;🎧 &lt;strong&gt;Hands-Free Audio Walk (ポケットの中のガイド)&lt;/strong&gt;: Real-time audio narration about trees, birds, local geography, and seasonal highlights without having to look down at your screen.&lt;/li&gt;
&lt;li&gt;🤖…&lt;/li&gt;
&lt;/ul&gt;&lt;/div&gt;
  &lt;/div&gt;
  &lt;div class="gh-btn-container"&gt;&lt;a class="gh-btn" href="https://github.com/jun-matsui/outscape" rel="noopener noreferrer"&gt;View on GitHub&lt;/a&gt;&lt;/div&gt;
&lt;/div&gt;


&lt;h2&gt;
  
  
  How I Built It
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;Frontend &amp;amp; Design System:&lt;/strong&gt; Modern semantic HTML5, Vanilla CSS with an Obsidian &amp;amp; Emerald glassmorphism design system, and Leaflet.js mapping with 100% free OpenStreetMap tiles.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Open-Weight AI Engine:&lt;/strong&gt; Configured around &lt;strong&gt;Google Gemma 2&lt;/strong&gt;, supporting &lt;strong&gt;WebGPU hardware acceleration&lt;/strong&gt;, custom Ollama/VM endpoints, and an embedded edge fallback.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Audio Synthesis:&lt;/strong&gt; Web Audio API pentatonic chimes + Web Speech API and ElevenLabs streaming for zero-latency, hands-free audio guidance.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Offline-First Resilience:&lt;/strong&gt; Outscape uses an embedded edge nature cache, ensuring that if you lose cell reception on a remote trail, the app continues to guide you without interruption.&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  Why Does Open Innovation Matter?
&lt;/h2&gt;

&lt;p&gt;When building an application for the outdoors, open-weight models and open innovation are far superior to closed proprietary APIs for three critical reasons:&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;
&lt;strong&gt;Connectivity in Nature:&lt;/strong&gt; Deep trails, mountains, and parks often have zero mobile signal. Closed APIs fail immediately when offline. Open-weight models can run directly on-device via WebGPU or local edge runtimes.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Location Privacy:&lt;/strong&gt; Hiking routes and daily walking habits are sensitive personal data. Outscape keeps all GPS telemetry on the client—never transmitted to third-party ad networks.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Zero-Cost Exploration:&lt;/strong&gt; Community trail groups, schools, and birdwatchers can use and adapt Outscape freely without paying per-token API charges.&lt;/li&gt;
&lt;/ol&gt;

&lt;h2&gt;
  
  
  Show Your Work: Agent Session Transcript
&lt;/h2&gt;

&lt;p&gt;To document the architectural decisions, prompt design, and engineering iterations behind Outscape, here is the curated agent session recorded with DevRelay:&lt;/p&gt;


&lt;div class="ltag-agent-session"&gt;
  &lt;div class="agent-session-header"&gt;
    
    &lt;span class="agent-session-tool-icon-badge" title="Gemini CLI"&gt;
&lt;/span&gt;
    &lt;span class="agent-session-title"&gt;Outscape: Engineering an Open-Weight AI Companion for the Outdoors&lt;/span&gt;
  &lt;/div&gt;

  &lt;div class="agent-session-scroll"&gt;

      &lt;div class="agent-session-message agent-session-user"&gt;
        &lt;div class="agent-session-role-badge agent-session-role-user"&gt;
          You
        &lt;/div&gt;
        &lt;div class="agent-session-content"&gt;
                &lt;div class="agent-session-text"&gt;
                  &lt;p&gt;Let's build Outscape for Hacktoberfest 2026 Touch Grass challenge: an open-source outdoor navigation and nature exploration companion that gets people off screens and into the wild.&lt;/p&gt;

                &lt;/div&gt;
        &lt;/div&gt;
      &lt;/div&gt;

      &lt;div class="agent-session-message agent-session-assistant"&gt;
        &lt;div class="agent-session-role-badge agent-session-role-assistant"&gt;
          Agent
        &lt;/div&gt;
        &lt;div class="agent-session-content"&gt;
                &lt;div class="agent-session-text"&gt;
                  &lt;p&gt;We will design Outscape around Pocket Mode: screen interaction is minimal. We leverage Google Gemma 2 open-weight architecture with WebGPU hardware acceleration for zero-telemetry local inference, ensuring it works even deep in backcountry trails with no signal.&lt;/p&gt;

                &lt;/div&gt;
              &lt;div class="agent-session-tool-call"&gt;
                
                  &lt;span class="agent-session-chevron"&gt;▸&lt;/span&gt;
                  &lt;span class="tool-name tool-name-viewfile"&gt;view_file&lt;/span&gt;
                    &lt;span class="tool-input-preview"&gt;index.html&lt;/span&gt;
                
                &lt;div class="agent-session-tool-detail"&gt;
                    &lt;div class="tool-section"&gt;
                      &lt;div class="tool-section-label"&gt;Input&lt;/div&gt;
                      &lt;pre class="agent-session-pre highlight html"&gt;&lt;code&gt;index.html&lt;/code&gt;&lt;/pre&gt;
                    &lt;/div&gt;
                    &lt;div class="tool-section"&gt;
                      &lt;div class="tool-section-label"&gt;Output&lt;/div&gt;
                      &lt;pre class="agent-session-pre highlight html"&gt;&lt;code&gt;Loaded Outscape PWA interface with WebGPU runtime and Leaflet trail mapping.&lt;/code&gt;&lt;/pre&gt;
                    &lt;/div&gt;
                &lt;/div&gt;
              &lt;/div&gt;
        &lt;/div&gt;
      &lt;/div&gt;

      &lt;div class="agent-session-message agent-session-user"&gt;
        &lt;div class="agent-session-role-badge agent-session-role-user"&gt;
          You
        &lt;/div&gt;
        &lt;div class="agent-session-content"&gt;
                &lt;div class="agent-session-text"&gt;
                  &lt;p&gt;How do we deliver hands-free trail guidance without requiring the user to look at the display?&lt;/p&gt;

                &lt;/div&gt;
        &lt;/div&gt;
      &lt;/div&gt;

      &lt;div class="agent-session-message agent-session-assistant"&gt;
        &lt;div class="agent-session-role-badge agent-session-role-assistant"&gt;
          Agent
        &lt;/div&gt;
        &lt;div class="agent-session-content"&gt;
                &lt;div class="agent-session-text"&gt;
                  &lt;p&gt;We implement ambient audio cues combined with contextual narration using Web Speech API and ElevenLabs audio streaming. As hikers encounter trail waypoints, local Gemma 2 models synthesize real-time nature observations and mindful 'Touch Grass' prompts directly to headphones.&lt;/p&gt;

                &lt;/div&gt;
        &lt;/div&gt;
      &lt;/div&gt;
  &lt;/div&gt;

  &lt;div class="agent-session-footer"&gt;
    &lt;span class="agent-session-meta"&gt;
        4 of 4 messages
    &lt;/span&gt;
  &lt;/div&gt;
&lt;/div&gt;


&lt;h2&gt;
  
  
  Prize Categories
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;strong&gt;Main Theme:&lt;/strong&gt; Touch Grass (Getting people off screens and into the natural world)&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Best Use of Gemma:&lt;/strong&gt; Utilizes Google Gemma 2 open-weight architecture with WebGPU hardware acceleration for contextual nature observations.&lt;/li&gt;
&lt;li&gt;
&lt;strong&gt;Best Use of ElevenLabs:&lt;/strong&gt; Implements dynamic voice streaming and narration triggers for hands-free audio walks.&lt;/li&gt;
&lt;/ul&gt;

</description>
      <category>devchallenge</category>
      <category>hf26challenge</category>
      <category>opensource</category>
      <category>ai</category>
    </item>
  </channel>
</rss>
