Legacy ETL replacements rarely fail because the new platform can't run the transformations. They fail because nobody can prove the new outputs are correct, or because a finance report three teams away breaks on the first Monday after cutover.
A legacy pipeline has usually been around long enough to collect undocumented fixes, silent type coercions, and consumers nobody remembers onboarding. Rewriting its SQL in a new tool is the easy part. Keeping the behavior that the business actually depends on is the hard part.
This article covers three concerns that decide whether a migration to Palantir Foundry goes smoothly:
- Incremental transformation: processing only what changed, instead of rebuilding everything nightly.
- Schema evolution: handling source and output schema changes without breaking consumers.
- Reconciliation: proving, with evidence, that the new pipeline produces equivalent results.
Running the legacy and Foundry pipelines side by side ties these together. Parallel operation turns cutover from a leap of faith into a decision backed by data.
Why Legacy ETL Migrations Are Difficult
Most legacy ETL estates share a familiar set of problems:
- Monolithic jobs. A single workflow extracts, cleans, joins, aggregates, and loads, so you can't migrate one piece without understanding all of it.
- Hidden dependencies. A downstream job reads an intermediate staging table that was never meant to be public.
- Hard-coded logic. Currency rates, region mappings, and exclusion lists live inside stored procedures.
- Fragile schedules. Job B starts at 02:30 because job A "usually finishes by then."
- Inconsistent schemas. The same customer ID is a string in one system and an integer in another.
- Limited observability. Failures surface as a stale dashboard, not as an alert.
- Locked-in consumers. Reports, extracts, and applications depend on exact table names, column orders, and even rounding behavior.
A big bang migration bets that every one of these has been found and handled before go-live. In practice, some will only be found after go-live, when the old system is gone and there's nothing to compare against.
A Safer Migration Pattern
An incremental approach reduces that risk by migrating in bounded slices. One common lifecycle:
- Inventory existing jobs, sources, targets, and schedules.
- Identify critical data products and who consumes them.
- Select a migration boundary small enough to validate fully, such as one subject area or one fact table and its dimensions.
- Build the corresponding Foundry pipeline.
- Run both pipelines in parallel on the same inputs.
- Compare outputs systematically.
- Resolve differences, classifying each as a new-pipeline bug, a legacy bug, or an accepted change.
- Shift consumers gradually.
- Retire the legacy pipeline once nothing reads from it.
This is an architectural pattern, not a mandatory Palantir methodology. Its value is that every step produces evidence, and the blast radius of a mistake stays limited to one slice.
Incremental Transforms in Foundry
Many legacy jobs follow a truncate-and-reload pattern: delete the target table, rebuild it from full source extracts. It's simple, but cost and runtime grow with total data volume rather than with the volume of change.
What Foundry provides
Foundry datasets are versioned through transactions. A build either commits a transaction or it doesn't, so a failed run doesn't leave a half-written output visible to consumers. Python transforms in Code Repositories support an @incremental decorator. When a transform runs incrementally, reading an input returns only the rows added since the last successful build, and the output can append rather than replace. Bumping the decorator's semantic_version forces a full recompute, which is useful when logic changes. Data Connection also supports incremental syncs for some source types, typically by tracking a monotonically increasing column such as an update timestamp or sequence ID. Pipeline Builder has its own incremental options; what's available can depend on your enrollment and pipeline type.
An important constraint: incremental transforms are built around appended data. If an upstream dataset is fully replaced by a snapshot, a downstream incremental transform generally can't run incrementally. Foundry datasets also aren't key-based tables with native upserts, so "latest state per key" has to be modeled explicitly.
Design concerns
-
Change detection. Prefer reliable change signals (CDC feeds, sequence numbers) over
updated_atcolumns that applications forget to maintain. - Partitioning. Partition outputs by a column consumers filter on, such as business date, so reprocessing stays targeted.
- Watermarks. Ingestion-side watermarks decide what gets pulled. Transform-side incrementality decides what gets processed. Keep the two concepts separate in your design.
- Late-arriving data. An order updated today for last month's date lands in today's increment. Any aggregate by business date must recompute the affected dates, not just append.
-
Idempotency. Transactions make commits atomic, but your logic must still be deterministic. Avoid
current_timestamp()inside business logic and deduplicate on stable keys.
Example: from full refresh to changelog plus current state
Legacy: a nightly job truncates fact_orders and reloads from a full ERP extract.
Redesign: ingest only changed rows into an append-only changelog, then derive the current state per order.
from pyspark.sql import functions as F, Window
from transforms.api import transform, incremental, Input, Output
@incremental(semantic_version=1)
@transform(
out=Output("/Sales/clean/orders_current"),
changes=Input("/Sales/clean/orders_changelog"),
)
def compute(changes, out):
new_rows = changes.dataframe() # only rows appended since last build
previous = out.dataframe("previous", schema=new_rows.schema)
w = Window.partitionBy("order_id").orderBy(
F.col("updated_at").desc(), F.col("ingested_at").desc()
)
latest = (
previous.unionByName(new_rows)
.withColumn("rn", F.row_number().over(w))
.filter("rn = 1")
.drop("rn")
)
out.set_mode("replace")
out.write_dataframe(latest)
The trade-off is visible: the extraction and changelog are incremental, but this current-state dataset is still rewritten each run. For very large tables, teams may partition the current-state output, or have consumers read from the changelog through a view. The decorator style above is the classic API; check your repository template for the syntax your Foundry version uses.
Managing Schema Evolution
Schema changes cause more migration incidents than transformation logic does, because they often fail silently.
| Change | Typical risk |
|---|---|
| Column added | Low, unless consumers use SELECT * or positional reads |
| Field renamed | Breaks consumers; can look like a dropped column plus a new null column |
| Type changed | Silent nulls on failed casts, precision loss, timezone shifts |
| Field removed | Breaks consumers and business rules that depend on it |
| Nested structure changed | Breaks flattening logic, often deep inside the pipeline |
| Source version upgrade | Several of the above at once, often undocumented |
Practices that help:
- Schema contracts. Define the canonical output schema explicitly and treat changes to it as reviewed changes.
- Explicit mappings. Map source fields to canonical fields by name, never by position.
- Version-aware transforms. Handle known source versions deliberately instead of hoping the shape stays stable.
-
Validation before promotion. Foundry Data Expectations let a transform fail a build when checks don't pass, for example
Check(E.primary_key("order_id"), "orders pk", on_error="FAIL"). Foundry also supports branching, so schema changes can be built and checked on a branch before merging. - Consumer impact analysis. Use Foundry's data lineage to see which downstream datasets and applications read a column before you change it.
Example: a source rename and type change
The ERP upgrade renames cust_id to customer_id and changes amount from string to decimal.
from pyspark.sql import functions as F, types as T
CANONICAL = {
"order_id": T.StringType(),
"customer_id": T.StringType(),
"amount": T.DecimalType(18, 2),
"updated_at": T.TimestampType(),
}
RENAMES = {"cust_id": "customer_id"} # known source v1 -> v2 mapping
def to_canonical(df):
for old, new in RENAMES.items():
if old in df.columns and new not in df.columns:
df = df.withColumnRenamed(old, new)
missing = [c for c in CANONICAL if c not in df.columns]
if missing:
raise ValueError(f"Schema contract violated, missing: {missing}")
return df.select([F.col(c).cast(t).alias(c) for c, t in CANONICAL.items()])
One catch: with Spark's default (non-ANSI) settings, a failed cast returns null instead of raising an error. Pair the mapping with a null-rate check on amount. Also, if the output schema of an incremental transform changes, appending to the existing output will typically conflict. Bumping semantic_version to force a full rebuild is the cleaner path.
Reconciliation Between Legacy and Foundry Pipelines
Matching row counts proves very little. Two tables can have identical counts while one has duplicated keys, swapped values, or nulls where the other has data.
Reconcile at several levels:
| Level | What it catches |
|---|---|
| Row counts | Gross load failures, filter differences |
| Key coverage | Records present in only one output |
| Duplicate detection | Join fan-out, broken dedup logic |
| Aggregate totals | Sums by date, region, product; rounding drift |
| Null rates per column | Failed casts, broken mappings |
| Record-level comparison | Value differences on matching keys |
| Business-rule validation | Rules that must hold regardless of either pipeline |
| Sample-based review | Human inspection where full comparison is impractical |
A common approach is to ingest the legacy output into Foundry via Data Connection, then compare in a dedicated reconciliation transform:
COMPARE_COLS = ["customer_id", "amount", "status"]
def row_hash(df):
normalized = [F.coalesce(F.trim(F.col(c).cast("string")), F.lit("<null>"))
for c in COMPARE_COLS]
return df.select("order_id", F.sha2(F.concat_ws("||", *normalized), 256).alias("h"))
legacy, foundry = row_hash(legacy_df).alias("l"), row_hash(foundry_df).alias("f")
result = legacy.join(foundry, "order_id", "full_outer").select(
"order_id",
F.when(F.col("f.h").isNull(), "only_legacy")
.when(F.col("l.h").isNull(), "only_foundry")
.when(F.col("l.h") != F.col("f.h"), "value_mismatch")
.otherwise("match").alias("status"),
)
Normalization rules matter more than the comparison itself: decimal scale, timestamp timezones, trailing whitespace, empty strings versus nulls. Agree on documented tolerances up front, and record every accepted difference with a reason. Some mismatches turn out to be legacy bugs, and the business has to decide whether to fix or preserve them.
Running Both Pipelines Without Downtime
"Without downtime" here means consumers keep getting correct data on schedule throughout the migration. Whether cutover itself is invisible depends on your architecture.
- Dual processing. Both pipelines run on the same inputs and schedule.
- Shadow validation. Foundry outputs are reconciled every run, but nobody consumes them yet.
- Consumer-by-consumer migration. Move low-risk consumers first, critical ones last.
- Read-path switching. Give consumers a stable interface, such as a warehouse view, so switching means redirecting the view, not editing every report. Where Foundry needs to feed an external system, Data Connection exports may handle this, depending on the target.
- Rollback planning. Define measurable rollback triggers, like mismatch rate above threshold or a missed freshness SLA, and keep the legacy path warm until the window closes.
- Monitoring. Foundry health checks can cover freshness, build status, and schema. Monitor the reconciliation results as well.
- Ownership. Name who fixes a broken build on each side during the transition. Ambiguity here causes most parallel-run incidents.
A Practical Migration Example
A hypothetical scenario. A retailer runs a nightly ETL job that pulls customers from a CRM, orders from an ERP, and payments from a payment gateway, then loads a reporting warehouse star schema used by dozens of reports.
Legacy pipeline. Full extracts land in staging, stored procedures build dimensions and fact_orders, and reports read the warehouse directly.
Foundry pipeline. Data Connection syncs land raw data, incremental where the source exposes a reliable change column. Python transforms build changelogs, apply the schema contracts, and produce current-state customer, order, and payment datasets. Data Expectations enforce primary keys and non-negative amounts.
Parallel validation. Both pipelines run nightly through at least one month-end close, since period-end adjustments and late payments often surface edge cases that daily runs don't.
Reconciliation. Early runs show value mismatches on amount. The cause is a legacy procedure that rounds at line level while the Foundry transform rounds at order level. Finance chooses line-level rounding, and the Foundry logic is adjusted to match.
Consumer migration. Operational dashboards move first through view redirection, then finance reports after a clean month-end.
Legacy retirement. Once query logs show no reads from the legacy tables for an agreed period, the jobs are disabled, archived, and eventually removed.
Common Migration Mistakes
- Migrating everything at once, which removes the ability to compare and roll back.
- Ignoring downstream dependencies, especially consumers of intermediate tables.
- Treating schema changes as column edits instead of contract changes with consumers.
- Validating only row counts.
- Switching consumers before reconciliation is stable across a full business cycle.
- No rollback criteria, so rollback becomes a debate during an incident.
- Unclear ownership during the parallel period.
- Retiring legacy jobs too early, before quarterly or annual processes have run.
Where Palantir Consulting Fits
Some organizations bring in Palantir consulting or implementation partners for specific parts of the program: migration architecture, legacy dependency analysis, incremental pipeline design, schema strategy, reconciliation frameworks, deployment planning, and operational handover. The most useful engagements leave the internal team able to own and extend the pipelines afterward. If you use outside help, make knowledge transfer an explicit deliverable.
Migration Checklist
- Source inventory: systems, extraction methods, volumes, change signals
- Dependencies: all consumers, including intermediate-table readers
- Transformation logic: documented, including undocumented fixes in legacy code
- Incremental strategy: append versus snapshot, late data, recompute triggers
- Schema compatibility: contracts, explicit mappings, versioning rules
- Data quality: expectations enforced at build time
- Reconciliation: multi-level checks with agreed tolerances
- Parallel execution: same inputs, same schedule, through a full business cycle
- Monitoring: freshness, build health, reconciliation trends
- Cutover: consumer order, read-path switching, communication plan
- Rollback: measurable triggers and a warm legacy path
- Legacy retirement: read-log evidence, archival, decommission date
Final Takeaway
Moving ETL to Palantir Foundry is not mainly a code-translation exercise. The real work is controlling correctness, dependencies, and operational risk while the business keeps running. Incremental processing keeps the new pipeline efficient, schema contracts keep it stable, and reconciliation gives you the evidence to cut over. Migrate in slices, compare relentlessly, and retire the old system only when the data says it's safe.
Top comments (0)