DEV Community

Cover image for Stop Choosing Between BigQuery and Databricks Based on Marketing Slides
Aniket Abhishek Soni
Aniket Abhishek Soni

Posted on

Stop Choosing Between BigQuery and Databricks Based on Marketing Slides

03:42 AM. The PagerDuty alert didn't wake me up—the sound of the laptop fan spinning up to jet-engine levels in my sleep-deprived brain did.

The myth that "BigQuery and Databricks are both just SQL engines, so they perform the same" is the kind of professional lie we tell ourselves to avoid reading the documentation. I was looking at a dashboard that usually loads in 12 seconds, now spinning for 14 minutes before throwing a RESOURCE_EXHAUSTED error. We were mid-migration, testing our core reporting layer against 4TB of Parquet files sitting on S3.

The team assumed that since we were using Databricks SQL (Serverless) and BQ’s BigLake, the abstraction layer would handle the heavy lifting. It didn't. We were trying to join a massive fact table—denormalized, partitioned by day, and sitting in open-format delta tables—against a dimension table that someone had accidentally stored as a thousand tiny JSON files.

What we saw

The symptom was simple: the query scheduler in Databricks was timing out, and BigQuery was racking up a $400 bill for a single exploratory analysis.

We fell for the classic trap: "The engine will figure it out." In our Databricks environment, the SQL Warehouse was set to Pro with Auto-stop enabled. When we threw a complex window function over a non-partitioned column, the spill-to-disk metrics hit the roof. We saw io.file.write.bytes spiking in the Ganglia metrics, but the throughput was abysmal.

On the BigQuery side, the INFORMATION_SCHEMA.JOBS_BY_PROJECT view showed us that we were scanning the entire dataset because our partitioning key (event_date) was buried inside a CASE statement. The engine wasn't pruning partitions because the optimizer couldn't prove the predicate was constant at planning time. We spent three hours chasing "network latency" and "cold starts," convinced the cloud provider was throttling us. They weren't. We were just writing garbage SQL that assumed the engine had infinite compute and perfect metadata.

Photo by luka Bruyninckx on Unsplash
Photo by luka Bruyninckx on Unsplash

Root cause

The failure wasn't the engine; it was the impedance mismatch between our file layout and the compute engine’s execution model.

In Databricks, we were hitting a classic Z-Order failure. Our date column was the partition key, but our queries were filtering by user_id. Without Z-Ordering the data on the user_id column, the engine performed a full table scan of every partition to find a single user's activity. The offending code looked like this:

-- The culprit: forcing a full scan because of the non-indexed filter
SELECT 
  user_id, 
  SUM(revenue) OVER(PARTITION BY user_id ORDER BY transaction_timestamp) as running_total
FROM delta.`s3://our-bucket/prod/transactions/`
WHERE user_id = '98273-a' -- This column is not Z-Ordered
Enter fullscreen mode Exit fullscreen mode

On the BigQuery side, the issue was our reliance on BigLake. We had configured the table to point to S3, but we hadn't defined the metadata cache properly. BigQuery was spending more time listing files in the S3 bucket via the storage API than actually executing the sum. It was a "Metadata Bottleneck." Every time we queried, BigQuery was hitting the S3 ListObjectsV2 API thousands of times. The cost wasn't just the compute; it was the overhead of the cloud provider’s storage interface being treated like a local SSD.

Photo by Markus Spiske on Unsplash
Photo by Markus Spiske on Unsplash

The fix

We stopped trying to make the engine "do it all" and started managing the physical data layout.

For Databricks, we ran an OPTIMIZE transactions ZORDER BY (user_id) command. This isn't a silver bullet—it's a maintenance job that costs money—but it immediately reduced the amount of data read by 90%. We moved the user_id to a clustered index equivalent within the Delta Lake storage layer.

For the BigQuery side, we forced the issue by enabling metadata_cache_mode = 'AUTOMATIC' and setting the metadata_cache_ttl to 60 minutes. We also had to switch from raw S3 URIs to a managed BigQuery external table definition that explicitly defined the partition layout.

The change in performance was binary. The query dropped from 14 minutes to 4 seconds. The cost dropped from $400 to $0.12. We weren't fighting a platform limitation; we were fighting our own laziness regarding metadata management.

What we changed so it never happens again

We implemented a strict "No Raw File Access" policy. If a data engineer wants to query raw S3 files, they use a staging layer. If it’s for production, it must be a managed table with defined partitioning and Z-Ordering.

We also introduced a cost-budgeting script that triggers a Slack alert if any single query exceeds $10 in BigQuery or 30 minutes of compute time in Databricks. We treat these alerts as P1 incidents. No "I'll optimize it later." If it hits the limit, it gets killed, and the developer has to explain the EXPLAIN plan in the pull request.

The biggest lesson? Don't trust the marketing slides that say "BigQuery is serverless" or "Databricks handles all the tuning." Both platforms are incredibly powerful, but they both expect you to know how they handle physical storage.

If you treat a cloud data warehouse like a local Postgres instance, you will get billed for your ignorance. If you treat it like a distributed system, it’s the most powerful tool you have. Stop writing SQL for a laptop and start writing SQL for the cluster. If you don't know where your data lives and how it's partitioned, the engine is just going to scan the entire internet and charge you for the privilege.

Cover photo by Albert Stoynov on Unsplash.

Top comments (0)