Workload
Performance tuning & optimization for Databricks → BigQuery
Spark tuning habits (shuffle optimization, partition overwrite, cluster sizing) don’t directly translate. We tune BigQuery layout, queries, and capacity so pruning works, bytes scanned stays stable, and refresh SLAs hold under real concurrency.
Quick answer
Spark tuning habits (shuffle optimization, partition overwrite, cluster sizing) don’t directly translate. We tune BigQuery layout, queries, and capacity so pruning works, bytes scanned stays stable, and refresh SLAs hold under real concurrency.
Back to pair pageContext
Why this breaks
Databricks performance is dominated by cluster runtime and Spark execution patterns (shuffles, partition sizing, file compaction). In BigQuery, cost and runtime are driven by bytes scanned, slot time, and how well queries prune partitions. After cutover, teams often keep Spark-era query shapes and lose pruning—leading to scan blowups, slow dashboards, and unpredictable spend. Common post-cutover symptoms:
- Queries scan entire facts because partition filters aren’t pushed down - Heavy joins reshuffle large datasets; BI refresh becomes slow and expensive - Incremental MERGE/apply touches too much history each run - Repeated parsing/extraction of semi-structured fields becomes expensive - Concurrency spikes cause slot contention and tail latency Optimization replaces “Spark tuning” with BigQuery-native pruning, layout, and governance—backed by baselines and regression gates.
Approach
How conversion works
- Baseline the top workloads: identify the most expensive and most business-critical queries/pipelines (dashboards, marts, MERGE/upserts). - Diagnose root causes: scan bytes, join patterns, partition pruning, repeated transforms, and apply scope. - Tune table layout: partitioning and clustering aligned to real access paths and refresh windows. - Rewrite for pruning and reuse: pruning-first filters, pre-aggregation/materializations, and typed extraction boundaries for semi-structured fields. - Capacity & cost governance: on-demand vs reservations posture, concurrency controls, and guardrails for expensive queries. - Regression gates: store baselines and enforce thresholds so improvements persist.
Coverage
Supported constructs
Representative tuning levers we apply for Databricks → BigQuery workloads.
| Source | Target | Notes |
|---|---|---|
| Spark shuffle-heavy joins | Pruning-first joins + layout alignment | Reduce scanned bytes and stabilize runtime under BI refresh. |
| Delta partition overwrite habits | Bounded apply windows + partition-scoped MERGE | Avoid touching more history than necessary each run. |
| Compaction/OPTIMIZE habits | Partitioning/clustering + materialization strategy | Keep pruning effective as data grows. |
| Semi-structured parsing in Spark | Typed extraction boundaries + reuse | Extract once, cast once, reuse everywhere. |
| Cluster sizing for performance | Slots/reservations + concurrency posture | Stabilize SLAs under peak usage and control spend. |
| Ad-hoc expensive queries | Governance guardrails + cost controls | Prevent scan blowups and surprise bills. |
Compare
How workload changes
| Topic | Databricks / Delta | BigQuery | Notes |
|---|---|---|---|
| Primary cost driver | Cluster runtime | Bytes scanned + slot time | Pruning and query shape dominate spend. |
| Tuning focus | Shuffle patterns, partition sizing, file layout | Partitioning/clustering + pruning-first SQL | Layout decisions become first-class levers. |
| Incremental apply | Partition overwrite and reprocessing windows common | Bounded MERGE/apply with explicit late windows | Correctness and cost depend on apply window design. |
| Concurrency planning | Cluster autoscaling and job scheduling | Slots/reservations + concurrency policies | Peak BI refresh requires explicit capacity posture. |
Examples
Examples
Illustrative BigQuery optimization patterns after Databricks migration: enforce pruning, pre-aggregate for BI, scope MERGEs, and store baselines for regression gates.
-- Pruning-first query shape (fact table partitioned by DATE(event_ts))
SELECT
country,
SUM(revenue) AS rev
FROM `proj.mart.fact_orders`
WHERE DATE(event_ts) BETWEEN @start_date AND @end_date
GROUP BY 1; -- Pre-aggregate pattern for repeated BI scans
CREATE OR REPLACE TABLE `proj.mart.orders_daily`
PARTITION BY d
CLUSTER BY country AS
SELECT
DATE(event_ts) AS d,
country,
SUM(revenue) AS rev,
COUNT(*) AS orders
FROM `proj.mart.fact_orders`
WHERE DATE(event_ts) BETWEEN @start_date AND @end_date
GROUP BY 1,2; -- Scope MERGE to affected partitions to avoid full-table scans
DECLARE min_d DATE DEFAULT (SELECT MIN(DATE(updated_at)) FROM `proj.ds.stg_orders`);
DECLARE max_d DATE DEFAULT (SELECT MAX(DATE(updated_at)) FROM `proj.ds.stg_orders`);
MERGE `proj.mart.fact_orders` t
USING `proj.ds.stg_orders` s
ON t.id = s.id
AND DATE(t.updated_at) BETWEEN min_d AND max_d
WHEN MATCHED THEN UPDATE SET t.status = s.status, t.amount = s.amount, t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (id, status, amount, updated_at) VALUES (s.id, s.status, s.amount, s.updated_at); -- Store performance baselines for regression gates (simplified)
CREATE TABLE IF NOT EXISTS `proj.control.perf_baselines` (
artifact_id STRING,
captured_at TIMESTAMP,
runtime_ms INT64,
bytes_scanned INT64,
slot_ms INT64
); Workload Assessment
Replace Spark tuning with BigQuery-native pruning
We identify your highest-cost migrated workloads, tune pruning and table layout, and deliver before/after evidence with regression thresholds—so performance improves and stays stable.
Book assessmentAvoid
Common pitfalls
- Defeating pruning: wrapping partition columns in functions/casts in WHERE clauses.
- Partitioning by the wrong key: choosing partitions that don’t match common filters and refresh windows.
- Clustering by folklore: clustering keys chosen without evidence from predicates and join keys.
- Unbounded MERGE/apply: incremental jobs touch too much history, causing scan bytes spikes.
- Over-materialization: too many intermediates without controlling refresh cost.
- Ignoring concurrency: BI refresh peaks overwhelm slots/reservations and create tail latency.
- No regression gates: performance improves once, then regresses silently after the next release.
Proof
Validation approach
- Baseline capture: runtime, bytes scanned, slot time, and output row counts for top queries/pipelines.
- Pruning checks: confirm partition pruning and predicate pushdown on representative parameters and boundary windows.
- Before/after evidence: demonstrate improvements in runtime and scan bytes; document tradeoffs.
- Correctness guardrails: golden queries and KPI aggregates ensure tuning doesn’t change semantics.
- Regression thresholds: define alerts (e.g., +25% bytes scanned or +30% runtime) and enforce via CI or scheduled checks.
- Operational monitors: dashboards for scan bytes, slot utilization, failures, and refresh SLA adherence.
Execution
Migration steps
A sequence that improves performance while protecting semantics.
-
01
Identify top cost and SLA drivers
Rank queries and pipelines by bytes scanned, slot time, and business criticality (dashboards, batch windows). Select a tuning backlog with clear owners.
-
02
Create baselines and targets
Capture current BigQuery job metrics (runtime, scan bytes, slot time) and define improvement targets. Freeze golden outputs so correctness doesn’t regress.
-
03
Tune layout: partitioning and clustering
Align partition keys to common filters and refresh windows. Choose clustering keys based on observed predicates and join keys—not guesses.
-
04
Rewrite for pruning and reuse
Apply pruning-aware SQL rewrites, reduce reshuffles, scope MERGEs/applies to affected partitions, and pre-aggregate where BI patterns repeatedly scan large facts.
-
05
Capacity posture and governance
Set reservations/on-demand posture, tune concurrency for BI refresh peaks, and implement guardrails to prevent scan blowups from new queries.
-
06
Add regression gates
Codify performance thresholds and alerting so future changes don’t reintroduce high scan bytes or missed SLAs. Monitor post-cutover metrics continuously.
FAQ
Frequently asked questions
Why did costs increase after moving from Databricks to BigQuery? +
Most often because pruning was lost: partition filters don’t translate cleanly, or filters defeat partition elimination. We tune layout and rewrite queries to reduce bytes scanned and stabilize spend.
How do you keep optimization from changing results? +
We gate tuning with correctness checks: golden queries, KPI aggregates, and edge-cohort diffs. Optimizations only ship when outputs remain within agreed tolerances.
Can you optimize MERGE/upsert pipelines too? +
Yes. We scope MERGEs to affected partitions, design staging boundaries, and validate performance with scan bytes/slot time baselines to prevent unpredictable runtime.
Do you cover reservations and concurrency planning? +
Yes. We recommend a capacity posture (on-demand vs reservations), concurrency controls for BI refresh spikes, and monitoring/guardrails so performance stays stable as usage grows.
Optimization Program
Prevent scan blowups with regression gates
Get an optimization backlog, tuned partitioning/clustering, and performance gates (runtime/bytes/slot thresholds) so future releases don’t reintroduce slow dashboards or high spend.
Next reads
Related pages
- Read more
End-to-end approach: what breaks, validation gates, and cutover plan.
- Read more
Migrate Delta pipelines to BigQuery with idempotency and late-data integrity gates.
- Read more
Convert Databricks SQL to BigQuery with semantic parity and golden-query validation.
- Read more
How we define parity contracts, thresholds, and evidence for sign-off.