Workload
Hadoop legacy ETL pipelines to BigQuery
Re-home Hadoop-era pipelines—partition overwrites, staging zones, UDF utilities, and orchestrated SQL chains—into BigQuery with an explicit run contract and validation gates that prevent drift and scan-cost surprises.
Quick answer
Re-home Hadoop-era pipelines—partition overwrites, staging zones, UDF utilities, and orchestrated SQL chains—into BigQuery with an explicit run contract and validation gates that prevent drift and scan-cost surprises.
Back to pair pageContext
Why this breaks
Hadoop ETL systems rarely fail because a query won’t translate. They fail because the run contract is implicit: partition overwrite conventions, file-based staging semantics, late-arrival reprocessing windows, and restart behavior encoded in Oozie coordinators and shell scripts. BigQuery can implement equivalent outcomes, but only if those hidden rules are extracted and enforced with deterministic patterns and validation gates. Common symptoms after cutover:
- Duplicate or missing data because overwrite/reprocessing semantics weren’t recreated - Late-arrival events are ignored (or double-counted) because window policy was implicit - Scan-cost spikes because partition predicates and staging boundaries stop pruning - Orchestration dependencies and retries change, turning failures into silent data issues - Schema drift from upstream feeds breaks typed targets without a defined policy
Approach
How conversion works
- Inventory & classify the estate: Hive/Impala/Spark SQL, staging zones, UDFs, and orchestrators (Oozie/Airflow/Cron). - Extract the run contract: keys, partition boundaries, watermarks, dedupe tie-breakers, late-arrival window policy, and failure/retry semantics. - Re-home ingestion and staging: landing tables + manifests, typed staging, and standardized audit columns. - Rebuild transforms using BigQuery-native patterns (landing → typed staging → dedupe → apply) with partitioning/clustering aligned to access paths. - Implement restartability: applied-window tracking, idempotency keys, deterministic ordering, and safe retries. - Re-home orchestration: Composer/Airflow or your runner with explicit DAG contracts, retries, alerts, and concurrency posture. - Gate cutover with evidence: golden outputs + incremental integrity simulations (reruns, backfills, late injections) and rollback-ready criteria.
Coverage
Supported constructs
Representative Hadoop-era ETL constructs we commonly migrate to BigQuery (exact coverage depends on your estate).
| Source | Target | Notes |
|---|---|---|
| Partition overwrite pipelines (dt partitions) | Partition-scoped apply (MERGE or overwrite-by-partition) | Preserve overwrite semantics without full-table rewrites. |
| Oozie coordinators / shell script chains | Composer/Airflow DAGs with explicit contracts | Dependencies, retries, and SLAs become first-class artifacts. |
| Hive/Impala staging zones (files) | Landing tables + typed staging | Replayable staging boundaries with audit columns and manifests. |
| Late-data reprocessing windows | Explicit late-arrival policy + staged re-apply | Behavior verified via late-injection simulations. |
| SCD Type-1 / Type-2 logic | MERGE + current-flag/end-date patterns | Backfills and late updates tested as first-class scenarios. |
| Hive schema drift | Typed staging + drift policy (widen/quarantine/reject) | Auditability for changing upstream payloads. |
Compare
How workload changes
| Topic | Hadoop legacy cluster | BigQuery | Notes |
|---|---|---|---|
| Incremental correctness | Often emerges from overwrite + coordinator conventions | Explicit windows, idempotency, and staged apply | Correctness becomes auditable and repeatable under retries/backfills. |
| Performance model | Avoid HDFS scans by partition predicates | Bytes scanned is the cost driver; pruning must be explicit | Staging boundaries and filters drive stable cost and runtime. |
| Orchestration | Oozie coordinators and script chains | Composer/Airflow/dbt DAGs with explicit contracts | Retries and alerts are modeled and monitored. |
| Schema evolution | Hive drift tolerated by downstream consumers | Typed staging with explicit drift policy | Prevents silent coercion and downstream surprises. |
Examples
Examples
Canonical BigQuery pattern for windowed loads: stage → dedupe deterministically → partition-scoped apply + applied-window tracking. Adjust keys, partitions, and casts to your model.
-- Applied-window tracking (restartability)
CREATE TABLE IF NOT EXISTS `proj.control.applied_windows` (
job_name STRING NOT NULL,
window_start DATE NOT NULL,
window_end DATE NOT NULL,
applied_at TIMESTAMP NOT NULL
); -- Stage → dedupe → apply (partition-scoped)
CREATE TEMP TABLE stg AS
SELECT
CAST(id AS STRING) AS id,
CAST(status AS STRING) AS status,
CAST(amount AS NUMERIC) AS amount,
event_ts,
SAFE_CAST(src_seq AS INT64) AS src_seq,
ingested_at
FROM `proj.raw.events`
WHERE event_date BETWEEN @start_d AND @end_d;
CREATE TEMP TABLE stg_dedup AS
SELECT *
FROM stg
QUALIFY ROW_NUMBER() OVER (
PARTITION BY id
ORDER BY event_ts DESC, src_seq DESC, ingested_at DESC
) = 1;
MERGE `proj.mart.fact_orders` t
USING stg_dedup s
ON t.id = s.id
AND DATE(t.updated_at) BETWEEN @start_d AND @end_d
WHEN MATCHED THEN UPDATE SET
t.status = s.status,
t.amount = s.amount,
t.updated_at = s.event_ts
WHEN NOT MATCHED THEN
INSERT (id, status, amount, updated_at)
VALUES (s.id, s.status, s.amount, s.event_ts); -- Mark window applied (idempotency marker)
INSERT INTO `proj.control.applied_windows` (job_name, window_start, window_end, applied_at)
VALUES (@job_name, @start_d, @end_d, CURRENT_TIMESTAMP()); Workload Assessment
Migrate legacy pipelines with the run contract intact
We inventory your Hadoop pipelines, formalize partition/late-data semantics, migrate a representative pipeline end-to-end, and produce parity evidence with cutover gates—without scan-cost surprises.
Book assessmentAvoid
Common pitfalls
- Assuming “SQL translation” equals pipeline migration: the hard part is operational semantics (windows, retries, late data).
- Partition semantics lost: overwrite-partition becomes append-only; duplicates appear.
- Dedupe instability: missing tie-breakers causes nondeterministic drift.
- Pruning defeated: filters wrap partition columns or cast in WHERE, causing scan bytes explosion.
- Unbounded applies: MERGEs or refreshes touch too much history each run.
- Schema drift surprises: upstream types widen/change; typed targets break without a drift policy.
- Orchestrator mismatch: coordinator-based dependencies aren’t mapped; freshness and correctness drift.
Proof
Validation approach
- Execution checks: pipelines run reliably under representative volumes and schedules.
- Structural parity: partition/window-level row counts and column profiles (null/min/max/distinct).
- KPI parity: aggregates by key dimensions for critical marts and dashboards.
- Incremental integrity (mandatory): - _Idempotency:_ rerun same window → no net change - _Restart simulation:_ fail mid-run → resume → correct final state - _Backfill safety:_ historical windows replay without drift - _Late-arrival:_ inject late corrections → only expected rows change - _Dedupe stability:_ duplicates eliminated consistently under retries
- Cost/performance gates: pruning verified; scan bytes/runtime thresholds set for top jobs.
- Operational readiness: retry/alerting tests, canary gates, and rollback criteria defined before cutover.
Execution
Migration steps
A sequence that keeps pipeline correctness measurable and cutover controlled.
-
01
Inventory pipelines, schedules, and dependencies
Extract pipeline graph: SQL jobs, staging zones, upstream feeds, and orchestrators (Oozie/Airflow/Cron). Identify business-critical marts and consumers.
-
02
Formalize the run contract
Define windows/high-water marks, business keys, deterministic ordering/tie-breakers, dedupe rules, late-arrival policy, and backfill boundaries. Make restartability explicit.
-
03
Rebuild transformations on BigQuery-native staging
Implement landing → typed staging → dedupe → apply with partitioning/clustering aligned to windows and access paths. Define schema evolution policy (widen/quarantine/reject).
-
04
Re-home orchestration and operations
Implement DAGs in Composer/Airflow/dbt: dependencies, retries, alerts, and concurrency. Recreate overwrite/reprocessing behavior with explicit late windows and idempotent markers.
-
05
Run parity and incremental integrity gates
Golden outputs + KPI aggregates, idempotency reruns, late-data injections, and backfill windows. Cut over only when thresholds pass and rollback criteria are defined.
FAQ
Frequently asked questions
Why is Hadoop ETL migration more than rewriting SQL? +
Because correctness lives in operational conventions: partition overwrite, reprocessing windows, retries, and script-based orchestration. We extract these rules into an explicit run contract and validate them with integrity gates.
How do you preserve overwrite-partition behavior in BigQuery? +
We implement partition-scoped apply (MERGE or overwrite-by-partition) with applied-window tracking and rerun simulations so outcomes remain stable under retries and backfills.
What about late-arriving data? +
We convert it into an explicit late-arrival window policy and staged re-apply strategy, then validate with late-injection simulations to prove only expected rows change.
How do you prevent BigQuery cost surprises? +
We design pruning-aware staging boundaries and choose partitioning/clustering aligned to windows. Validation includes scan bytes/runtime baselines and regression thresholds for your top jobs.
Migration Acceleration
Cut over pipelines with proof-backed gates
Get an actionable migration plan with incremental integrity tests (reruns, late data, backfills), reconciliation evidence, and cost/performance baselines—so pipeline cutover is controlled and dispute-proof.
Next reads
Related pages
- Read more
End-to-end approach: what breaks, validation gates, and cutover plan.
- Read more
How we define parity contracts, thresholds, and evidence for sign-off.
- Read more
Keep spend predictable and catch regressions early.