Prophecy low-code data engineering is a visual, drag-and-drop pipeline builder that does something most low-code tools do not: instead of running your work inside a proprietary black box, it compiles the canvas to open-source code — readable PySpark or Scala Spark for data pipelines, or dbt-style SQL models for warehouse pipelines — and commits that code to your own Git repository. You wire a few boxes together on a graph, and Prophecy generates the exact Spark job a senior engineer would have hand-written, tests and all.
That is a genuinely different shape from the two options enterprise data teams reached for before it: a legacy visual ETL suite (Informatica, Alteryx, Ab Initio, SSIS) whose logic is trapped in an opaque metadata format only that vendor can execute, or a purely hand-coded Spark stack that only your most senior engineers can safely touch. Prophecy sits in the middle — visual enough that analysts and junior engineers are productive, but code-first enough that the output is standard, version-controlled, runnable-without-the-vendor Spark and SQL. This guide walks through the five ideas an interviewer or an architecture-review board will actually probe — the gem-and-pipeline model, the visual↔code round-trip on Git, execution on Spark and Databricks via Fabrics, and reusable subgraphs with tests and column-level lineage — and pairs each with a Solution-Tail interview answer: code, a step-by-step trace, an output table, then a concept-by-concept breakdown of why it works.
When you want hands-on reps immediately after reading, drill the Spark SQL practice library →, rehearse the extract-and-load shapes on the ETL practice set →, and sharpen the transform logic Prophecy generates on the data-transformation practice set →.
On this page
- Why Prophecy changes enterprise pipeline building in 2026
- Gems, pipelines & the visual canvas
- The visual ↔ code round-trip on Git
- Running on Spark & Databricks with Fabrics
- Reusable subgraphs, tests & column-level lineage
- Cheat sheet — Prophecy recipes
- Frequently asked questions
- Practice on PipeCode
1. Why Prophecy changes enterprise pipeline building in 2026
Prophecy is a code generator wearing a visual UI — that one fact decides where it fits
The one-sentence invariant: Prophecy is not a runtime you send data through; it is a visual editor that emits ordinary Spark and SQL code, which then runs on your own engine. Everything that makes Prophecy defensible to a platform team follows from that. There is no proprietary execution engine to license per row, no metadata format only the vendor can read, and no "export to a dead PDF of a diagram" migration cliff — the artifact Prophecy produces is a Git repo full of PySpark files or SQL models you could keep running if Prophecy vanished tomorrow.
The low-code-without-lock-in split — what Prophecy does and deliberately does not do.
- Authoring. Prophecy gives you a drag-and-drop canvas of gems (transform nodes) that snap into a pipeline graph, with per-gem data previews and a schema that flows through the graph as you build.
- Code generation. Every visual pipeline compiles to standard code: PySpark or Scala for Spark projects, dbt-style SQL models for SQL projects. That code is the source of truth and lives in Git.
- Execution is out of scope on purpose. Prophecy does not run your data on its own cluster; the generated Spark runs on Databricks (or any Spark), and the generated SQL runs on your warehouse via dbt. Prophecy stays a design-and-generate layer, which is exactly why the output has no lock-in.
Where Prophecy sits against the alternatives.
- vs hand-written Spark. A raw PySpark codebase is maximally flexible but gates every change behind a senior engineer and a code review. Prophecy lets an analyst build the same pipeline visually while still producing the same reviewable PySpark — you widen the contributor pool without lowering the artifact quality.
- vs legacy Informatica / Alteryx / Ab Initio. Those visual suites also let non-experts build pipelines, but the logic is trapped in a vendor metadata format and runs only on the vendor engine. Prophecy's visual output is open-source Spark/SQL in your repo, so there is no proprietary runtime to keep paying for and no black box to reverse-engineer during an audit.
- vs Fivetran / Airbyte. Managed connectors solve extract-load for sources that already have a connector. Prophecy solves transform — the joins, aggregations, and business logic downstream of ingestion — and generates the Spark/SQL that expresses it.
What interviewers listen for.
- Do you say "Prophecy compiles the canvas to code, it is not a runtime" in the first sentence? — senior signal.
- Do you frame the value as "low-code authoring, code-first artifact" rather than "no-code magic"? — required framing.
- Do you reach for Prophecy when the goal is "migrate Informatica to Spark without hand-rewriting 4,000 mappings" or "let analysts contribute reviewable Spark", not as a Fivetran replacement? — senior signal.
- Do you mention Git, tests, and column-level lineage as built-ins, not afterthoughts? — the enterprise point.
Worked example — three boxes that become a real Spark job
Detailed explanation. The canonical Prophecy "hello world" is a three-gem pipeline — a Source that reads a table, a Reformat that derives a column, and a Target that writes the result. It looks like a toy flowchart, and that is the point: the same three boxes that transform a five-row table generate the same PySpark structure that a production job would, because Prophecy always emits one function per gem and a main that wires them together. No notebook, no boilerplate SparkSession glue that you maintain by hand.
Question. Read an orders table, add a total = qty * price column, and write to orders_enriched. Show the shape of the Spark code Prophecy generates from the three gems.
Input.
| order_id | qty | price |
|---|---|---|
| 1 | 2 | 10.0 |
| 2 | 3 | 4.0 |
Code.
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql import functions as F
def Source(spark: SparkSession) -> DataFrame:
return spark.read.table("sales.orders")
def Reformat(spark: SparkSession, in0: DataFrame) -> DataFrame:
return in0.withColumn("total", F.col("qty") * F.col("price"))
def Target(spark: SparkSession, in0: DataFrame) -> None:
in0.write.mode("overwrite").saveAsTable("sales.orders_enriched")
def pipeline(spark: SparkSession) -> None:
df_src = Source(spark)
df_ref = Reformat(spark, df_src)
Target(spark, df_ref)
Step-by-step explanation. The Source gem becomes a function that returns a DataFrame from sales.orders; it declares where data comes from and nothing else. The Reformat gem becomes a pure DataFrame → DataFrame function that adds the total column with withColumn. The Target gem becomes a sink function that writes the table. The pipeline driver wires them in dependency order — exactly the DAG you drew on the canvas — so the visual graph and the code are the same object viewed two ways.
Output.
| order_id | qty | price | total |
|---|---|---|---|
| 1 | 2 | 10.0 | 20.0 |
| 2 | 3 | 4.0 | 12.0 |
Rule of thumb. If you can draw the pipeline as boxes and arrows, Prophecy can generate it as functions and a driver — the toy example and the 200-gem production pipeline differ only in how many gems there are, never in the code shape.
2. Gems, pipelines & the visual canvas
A gem is one transform node, a pipeline is a DAG of gems, and each gem compiles to a function
Prophecy has one authoring primitive you compose endlessly — the gem — and an interviewer who asks "how is a Prophecy pipeline structured?" wants the gem-to-pipeline-to-code chain in order. Get this vocabulary crisp and the whole tool snaps into focus.
The core objects.
- Gem — one transform node. A gem is a single box on the canvas that performs one operation. Built-in gems cover the standard Spark verbs: Source and Target (read/write), Reformat (select/derive columns), Aggregate (group-by), Join, Filter, OrderBy, Deduplicate, SetOperation (union/intersect), and FlattenSchema for nested data. Each gem carries its own configuration and its own output schema.
-
Pipeline — a DAG of gems. Wiring gem outputs to gem inputs builds a directed acyclic graph. Data flows as typed Spark
DataFrames along the edges; the schema at each edge is known at design time, so Prophecy can validate column references before you ever run. - Project — a Git-backed collection of pipelines. A project bundles pipelines, subgraphs, tests, and generated code into one repository. It is the unit that maps to a Git repo.
What every gem gives you.
-
A generated function. Each gem compiles to exactly one function in the project code —
Reformat(spark, in0),Aggregate(spark, in0), and so on — so the code is as modular as the diagram. -
A visual expression builder. Inside a gem you write column expressions in SQL-like syntax (
qty * price,upper(name)), and Prophecy translates them into the target dialect (Sparkfunctionscalls for PySpark, raw SQL for SQL projects). -
An interactive preview. You can run the pipeline up to any gem and see a sample of the
DataFrameat that point, which is the low-code equivalent of a breakpoint.
Why the gem is the powerful unit.
- Because each gem is an isolated function, a pipeline is trivially testable gem-by-gem and readable diff-by-diff in a pull request.
- Because gems carry schema, Prophecy catches "column does not exist" at design time, not at 3 a.m. in production.
- Because the gem set is extensible, a platform team can ship custom gems (via the Gem Builder) that encode company standards, and every analyst gets them as drag-and-drop boxes.
Worked example — a Filter gem and a Reformat gem in one pipeline
Detailed explanation. Real pipelines chain several gems. Here a two-gem transform keeps only paid orders (a Filter gem) and then standardizes the customer name (a Reformat gem). Each gem becomes its own function, and the pipeline driver composes them — the canvas order is the call order.
Question. Build a pipeline that filters orders to status = 'paid' and uppercases customer, and show the two generated gem functions plus the driver.
Input. Four orders, two of them paid, returned by the Source gem.
Code.
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql import functions as F
def FilterPaid(spark: SparkSession, in0: DataFrame) -> DataFrame:
return in0.filter(F.col("status") == "paid")
def StandardizeCustomer(spark: SparkSession, in0: DataFrame) -> DataFrame:
return in0.withColumn("customer", F.upper(F.col("customer")))
def pipeline(spark: SparkSession) -> None:
src = spark.read.table("sales.orders")
paid = FilterPaid(spark, src)
clean = StandardizeCustomer(spark, paid)
clean.write.mode("overwrite").saveAsTable("sales.orders_clean")
Step-by-step explanation. The Filter gem becomes FilterPaid, a DataFrame → DataFrame function applying one predicate. The Reformat gem becomes StandardizeCustomer, applying one withColumn. The driver reads the source, threads it through the two gem functions in canvas order, and writes the sink. Because each gem is a separate function, a reviewer reading the pull request sees exactly two logical changes, cleanly named.
Output.
| order_id | customer | status |
|---|---|---|
| 1 | ADA | paid |
| 3 | LINUS | paid |
Rule of thumb. One gem = one function = one reviewable unit. Keep each gem doing a single verb and your generated code stays as clean as your diagram.
Prophecy interview question on the gem/pipeline model
Question. An interviewer describes a Reformat gem that should derive three columns — total = qty * price, year = year(order_ts), and a normalized email = lower(trim(email)) — from an orders DataFrame. Write the single Spark transform this gem generates, keeping it one projection so the physical plan stays flat.
Solution Using a single projection with column expressions
Code.
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql import functions as F
def Reformat(spark: SparkSession, in0: DataFrame) -> DataFrame:
return in0.select(
"order_id",
"qty",
"price",
(F.col("qty") * F.col("price")).alias("total"),
F.year(F.col("order_ts")).alias("year"),
F.lower(F.trim(F.col("email"))).alias("email"),
)
Step-by-step trace.
| step | expression applied | column produced |
|---|---|---|
| 1 | qty * price |
total |
| 2 | year(order_ts) |
year |
| 3 | lower(trim(email)) |
email (overwritten) |
| 4 | one select projection |
flat plan |
- A Reformat gem is a projection: it maps input columns to output columns in a single
select, so all derivations land in one plan node. -
totalmultiplies two columns;yearextracts the year from a timestamp;emailis normalized by nestingloweraroundtrim. - Because everything is one
select, Spark's Catalyst optimizer fuses the expressions into a single project stage — no extra shuffle or scan. - Naming each derived column with
.alias(...)is what makes the gem's output schema explicit and lets the next gem referencetotalandyearby name.
Output:
| order_id | qty | price | total | year | |
|---|---|---|---|---|---|
| 1 | 2 | 10.0 | 20.0 | 2026 | ada@x.io |
Why this works — concept by concept:
-
Projection gem — a Reformat compiles to one
select, so N derived columns cost one plan node, not N chainedwithColumncalls that stack projections. -
Column expressions — writing
qty * pricein the gem and letting Prophecy emitF.col("qty") * F.col("price")keeps the visual and the code semantically identical. - Explicit aliases — every derived column is named, so the downstream schema is deterministic and design-time validation can catch typos.
-
Optimizer-friendly — a single flat projection lets Catalyst fuse expressions, which is why "one gem, one select" is both cleaner and faster than a pile of
withColumns. - Cost — the transform is O(rows) with no shuffle; a projection is a narrow transformation, so it adds no stage boundary.
Spark SQL
Topic — spark-sql
Spark SQL transform and projection problems
3. The visual ↔ code round-trip on Git
Every visual edit regenerates readable code, and the whole project lives in your Git repo
The feature that sells Prophecy to a platform team is the round-trip: the canvas and the code are two views of the same artifact, kept in sync bidirectionally, and both are stored in ordinary version control. Edit a gem and the code regenerates; edit the code and the canvas re-parses. This is the difference between a low-code tool you cannot audit and one that behaves like a normal software repository.
How the round-trip works.
- Visual → code. Every change on the canvas — adding a gem, editing an expression, rewiring an edge — immediately regenerates the corresponding function in the project's source files. There is no separate "export" step; the code is the save format.
-
Code → visual. Because Prophecy parses the generated code back into the graph, an engineer can open
pipeline.pyin their IDE, edit it, push, and see the canvas reflect the change. The code and diagram cannot drift apart. - Standard, readable output. The generated PySpark/SQL is not obfuscated machine output; it reads like code a person wrote, with named functions per gem, so it survives a human code review.
Git is the backbone, not a bolt-on.
- Projects are repositories. A Prophecy project maps to a Git repo (GitHub, GitLab, Bitbucket, Azure DevOps, or Prophecy-managed Git). Pipelines, subgraphs, tests, and configs are all files.
- Branch, commit, PR, merge. You develop on a branch, commit visual changes as normal diffs, open a pull request, get a review, and merge — the full software lifecycle, driven from a visual tool.
- CI/CD like any codebase. Because the artifact is code, you plug the repo into your existing CI to run tests and into your existing CD to deploy, with no Prophecy-specific pipeline server in the critical path.
Why "no lock-in" is a technical claim, not a slogan.
- The generated Spark job runs on any Spark cluster with
spark-submit; the generated SQL runs through dbt. Neither needs Prophecy at execution time. - If you stopped using Prophecy, you keep a working, readable, tested Spark/SQL codebase — the exit cost is "you lose the visual editor," not "your pipelines stop running."
Worked example — an Aggregate gem and its Git diff
Detailed explanation. The clearest way to see the round-trip is to add one Aggregate gem and read the diff it produces. The gem groups orders by customer and sums the total; the generated function is a plain groupBy(...).agg(...), and the Git diff a reviewer sees is a single new named function — reviewable in seconds.
Question. Add an Aggregate gem that computes revenue = sum(total) and orders = count(*) per customer, and show the generated function that lands in the commit.
Input.
| customer | total |
|---|---|
| ada | 20.0 |
| ada | 5.0 |
| linus | 12.0 |
Code.
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql import functions as F
def AggregateByCustomer(spark: SparkSession, in0: DataFrame) -> DataFrame:
return (
in0.groupBy("customer")
.agg(
F.sum("total").alias("revenue"),
F.count(F.lit(1)).alias("orders"),
)
)
Step-by-step explanation. The Aggregate gem's group-by key (customer) becomes groupBy("customer"), and each aggregate expression becomes one .agg(...) entry with an alias. Prophecy writes this as a single new function; the commit diff is + def AggregateByCustomer(...) plus one wiring line in the driver. A reviewer reads the diff exactly as they would review a hand-written Spark change — because it is one.
Output.
| customer | revenue | orders |
|---|---|---|
| ada | 25.0 | 2 |
| linus | 12.0 | 1 |
Rule of thumb. If a visual change does not produce a clean, reviewable code diff, you are drawing something the code cannot express cleanly — split it into more gems until each diff is one idea.
Prophecy interview question on aggregation transforms
Question. An interviewer asks you to compute, per customer, the total revenue and the number of distinct order days, then keep only customers with revenue over 100. Write the Spark aggregation this gem would generate.
Solution Using groupBy with a filtered aggregate
Code.
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql import functions as F
def CustomerRollup(spark: SparkSession, in0: DataFrame) -> DataFrame:
return (
in0.groupBy("customer")
.agg(
F.sum("total").alias("revenue"),
F.countDistinct(F.to_date("order_ts")).alias("active_days"),
)
.filter(F.col("revenue") > 100)
)
Step-by-step trace.
| customer | rows in | sum(total) | distinct days | kept? |
|---|---|---|---|---|
| ada | 3 | 140.0 | 2 | yes |
| linus | 2 | 60.0 | 2 | no |
| grace | 4 | 220.0 | 3 | yes |
-
groupBy("customer")collapses each customer's rows into one group, triggering one shuffle keyed oncustomer. -
sum("total")accumulates revenue andcountDistinct(to_date(order_ts))counts unique calendar days, both computed in the same aggregation pass. - The post-aggregation
filter(revenue > 100)is a HAVING clause: it runs on the grouped rows, droppinglinuswhile keepingadaandgrace. - Ordering matters — filtering after the aggregate means the predicate sees
revenue, a column that only exists post-group; a Filter gem placed before the Aggregate could not reference it.
Output:
| customer | revenue | active_days |
|---|---|---|
| ada | 140.0 | 2 |
| grace | 220.0 | 3 |
Why this works — concept by concept:
-
Group-by shuffle — the aggregation repartitions rows by
customerso all of one customer's rows meet on one executor; this is the single shuffle the whole transform costs. -
Multiple aggregates, one pass —
sumandcountDistinctare computed together, so two metrics cost one scan and one shuffle, not two. -
HAVING via post-filter — filtering after
aggis the Spark equivalent of SQLHAVING; the predicate references the aggregatedrevenue, which is only defined after the group. -
countDistinct semantics — counting distinct
to_date(order_ts)answers "how many active days," a different question fromcount(*); picking the right counter is the correctness crux. -
Cost — O(rows) to scan plus one O(rows) shuffle keyed on
customer; the final filter is O(groups), typically tiny.
ETL
Topic — etl
ETL aggregation and rollup problems
4. Running on Spark & Databricks with Fabrics
A Fabric binds a pipeline to where it runs — Spark on a cluster, or SQL on a warehouse
Authoring produces code; something has to run it. In Prophecy that something is a Fabric — the named execution environment that holds the connection, credentials, and cluster or warehouse settings your pipeline runs against. The same canvas can point at different Fabrics (dev, staging, prod), and the project language decides what the generated code targets: Spark projects emit PySpark/Scala that runs on a cluster, SQL projects emit dbt models that run on a warehouse.
What a Fabric is.
- The execution environment. A Fabric bundles where and how a pipeline runs: a Databricks workspace and cluster, a Spark-on-EMR/Dataproc/Livy endpoint, or a SQL warehouse (Databricks SQL, Snowflake, BigQuery) for SQL projects.
- Connection + credentials + compute. It carries the workspace URL, auth token/secret, and cluster size, so switching from a small dev cluster to a large prod cluster is a Fabric swap, not a code change.
- Interactive and scheduled. During development the Fabric backs the live per-gem preview; for production you deploy the pipeline as a scheduled Databricks Job or an Airflow DAG that calls the same generated code.
Spark projects vs SQL projects.
-
Spark project → PySpark/Scala. Gems compile to
DataFrametransformations; the job is submitted to a Spark cluster. Choose this for large-scale distributed transforms, complex joins, and ML feature pipelines. -
SQL project → dbt. Gems compile to SQL models with
ref()/source()dependencies that run in-warehouse via dbt. Choose this when the data already lives in the warehouse and you want warehouse-native, dbt-tested models. - Same authoring, different target. The gem canvas looks the same either way; the project language decides the dialect and the engine, which is how one tool serves both Spark and SQL teams.
What interviewers probe about execution.
- Development ≠ deployment. Interactive runs use an attached cluster for fast feedback; production runs are scheduled jobs on right-sized compute. Conflating the two is a classic junior mistake.
-
Cluster sizing is a Fabric concern. Spilling, out-of-memory, and shuffle-partition tuning are properties of the Fabric's compute, not the visual pipeline — the generated Spark obeys the same
spark.sql.shuffle.partitionsyou would tune by hand.
Worked example — the same join as Spark and as SQL
Detailed explanation. Because the project language chooses the dialect, the same visual Join gem generates PySpark in a Spark project and SQL in a SQL project. Seeing both makes the "same canvas, different engine" point concrete: an orders-to-customers inner join is one gem, two outputs.
Question. Join orders to customers on customer_id to attach region. Show what the Join gem generates in a Spark project versus a SQL project.
Input.
| order_id | customer_id |
|---|---|
| 1 | 7 |
| 2 | 8 |
Code.
## Spark project (PySpark) — generated Join gem
def JoinCustomers(spark, orders, customers):
return orders.join(customers, on="customer_id", how="inner") \
.select("order_id", "customer_id", "region")
-- SQL project (dbt model) — generated Join gem
select o.order_id, o.customer_id, c.region
from {{ ref('orders') }} o
join {{ ref('customers') }} c
on o.customer_id = c.customer_id
Step-by-step explanation. In the Spark project the Join gem emits orders.join(customers, on="customer_id", how="inner") and a projection selecting region. In the SQL project the identical gem emits an ANSI JOIN with dbt ref() calls that wire model dependencies. The canvas, the join key, and the join type are the same; only the target dialect and engine differ, decided by the project language and the Fabric.
Output.
| order_id | customer_id | region |
|---|---|---|
| 1 | 7 | EU |
| 2 | 8 | US |
Rule of thumb. Pick the project language by where the data lives and how big it is — Spark for distributed scale, SQL/dbt for warehouse-native models — then let the same gem canvas target either.
Prophecy interview question on deduplicating with Spark SQL
Question. Orders arrive with duplicate rows per order_id (retries), each with an updated_at. Keep exactly one row per order_id — the latest updated_at — using Spark SQL, the way a Deduplicate gem would generate it. How do you make it deterministic?
Solution Using ROW_NUMBER over a partition
Code.
with ranked as (
select
*,
row_number() over (
partition by order_id
order by updated_at desc, ingest_seq desc
) as rn
from orders
)
select * except (rn)
from ranked
where rn = 1
Step-by-step trace.
| order_id | updated_at | ingest_seq | rn | kept |
|---|---|---|---|---|
| 7 | 12:30 | 2 | 1 | yes |
| 7 | 12:30 | 1 | 2 | no |
| 8 | 11:00 | 1 | 1 | yes |
| 8 | 10:00 | 1 | 2 | no |
-
partition by order_idgroups the duplicate rows for each order together so ranking is per-key. -
order by updated_at descputs the freshest row first; the tie-breakingest_seq descmakes the winner deterministic when two rows share the sameupdated_at. -
row_number()assigns 1 to the winner of each partition, 2..N to the rest. -
where rn = 1keeps exactly one row perorder_id, andexcept (rn)drops the helper column so the output schema matches the input.
Output:
| order_id | updated_at | ingest_seq |
|---|---|---|
| 7 | 12:30 | 2 |
| 8 | 11:00 | 1 |
Why this works — concept by concept:
-
Window partition —
partition by order_idturns "dedupe per key" into a ranking problem solved inside each partition, no self-join required. -
Deterministic tie-break — adding
ingest_seq descto theorder byremoves the nondeterminism of ties, so re-running yields the same survivor every time. -
row_number vs rank —
row_numberguarantees a single winner (no ties at rank 1), which is exactly the "keep one" requirement;rankcould keep several. - except (rn) — projecting away the helper column keeps the deduplicated output schema-identical to the source, so downstream gems are unaffected.
-
Cost — one shuffle to partition by
order_idplus an in-partition sort, O(rows log rows) per partition; a Deduplicate gem generates exactly this plan.
Spark SQL
Topic — spark-sql
Window-function and dedup problems in Spark SQL
5. Reusable subgraphs, tests & column-level lineage
Package gems into reusable subgraphs, unit-test them, and trace every column
The three features that make Prophecy an enterprise tool rather than a personal one are reuse, testing, and lineage. A reusable subgraph turns a proven set of gems into a shareable component; gem-level unit tests pin behaviour in CI; and column-level lineage answers "if I change this field, what breaks?" across every pipeline in the org. These are the capabilities a legacy visual ETL suite either lacks or hides behind an opaque metadata store.
Reusable subgraphs.
- A subgraph is a group of gems as one node. You select several gems — say, the five that standardize a customer record — and package them into a reusable subgraph with defined input and output ports. It appears on the canvas as a single box.
- Build once, reuse everywhere. Every pipeline that needs "standardize customer" drops in the same subgraph; a fix to the subgraph propagates to all consumers, so business logic is DRY instead of copy-pasted across 40 pipelines.
- Parameterizable. Subgraphs (and custom gems built with the Gem Builder) take configuration, so one component adapts to slightly different inputs without forking.
Tests that run in CI.
- Gem-level unit tests. You attach sample input rows and expected output rows to a gem; Prophecy generates a real unit test (pytest for PySpark, dbt tests for SQL) that runs in your CI on every pull request.
- Fail the build, not production. Because tests are code in the repo, a transform regression fails the PR check — the same guardrail a hand-written Spark codebase has, now automatic from the visual definition.
- Data-quality expectations. Beyond unit tests, expectation checks (not-null, uniqueness, accepted ranges) can gate a load so bad data fails loudly instead of corrupting a table.
Column-level lineage.
-
Field-to-field tracing. Prophecy captures how each output column was derived from input columns across gems and pipelines, so you can trace
revenueback through every transform to its source columns. - Impact analysis. Before renaming or dropping a source column, lineage shows every downstream pipeline and report that depends on it — the difference between a safe change and a Monday-morning outage.
Worked example — a reusable "standardize customer" subgraph
Detailed explanation. The everyday reuse pattern is a subgraph that encapsulates a few cleaning gems behind one input and one output port. Its generated code is a function that takes a DataFrame and returns a DataFrame, so any pipeline can call it — the visual box and the reusable function are, again, the same thing.
Question. Package "trim + lowercase email, uppercase name, drop test rows" into a reusable subgraph and show the generated function that pipelines call.
Input. A raw customers DataFrame with mixed-case and test rows.
Code.
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql import functions as F
def StandardizeCustomer(spark: SparkSession, in0: DataFrame) -> DataFrame:
return (
in0.withColumn("email", F.lower(F.trim(F.col("email"))))
.withColumn("name", F.upper(F.col("name")))
.filter(~F.col("email").endswith("@test.local"))
)
Step-by-step explanation. The subgraph's three gems become three chained transformations inside one function: normalize email, uppercase name, and drop internal test accounts. Because it is exposed as a single DataFrame → DataFrame function with one input port, every pipeline reuses it by wiring one box — and a fix here updates all of them at once.
Output.
| name | |
|---|---|
| ADA | ada@x.io |
| GRACE | grace@y.io |
Rule of thumb. If the same three-or-more gems appear in two pipelines, promote them to a reusable subgraph — you trade a one-time packaging step for org-wide consistency and a single place to fix bugs.
Prophecy interview question on testing a transform
Question. You must guarantee that the "standardize customer" transform always lowercases email and never lets a @test.local row through. Write the unit test a Prophecy gem test generates so this is enforced in CI, and explain what a failure would catch.
Solution Using a gem unit test with expected output
Code.
from pyspark.sql import SparkSession
from mypipeline.gems import StandardizeCustomer
def test_standardize_customer(spark: SparkSession):
src = spark.createDataFrame(
[("Ada", " ADA@X.IO "), ("Bot", "bot@test.local")],
["name", "email"],
)
out = StandardizeCustomer(spark, src)
rows = {r["name"]: r["email"] for r in out.collect()}
assert rows == {"ADA": "ada@x.io"} # Bot dropped, Ada normalized
assert out.count() == 1 # test row filtered out
Step-by-step trace.
| input (name, email) | after normalize | after filter | in expected? |
|---|---|---|---|
| Ada, " ADA@X.IO " | ADA, ada@x.io | kept | yes |
| Bot, bot@test.local | BOT, bot@test.local | dropped | n/a |
- The test builds a tiny
DataFramewith one good row and one@test.localrow, then calls the exact gem function under test. -
StandardizeCustomertrims and lowercases the email, uppercases the name, and filters the@test.localrow out. - The assertion pins the whole expected output — one row,
{"ADA": "ada@x.io"}— so any regression (a missinglower, a broken filter) fails the assertion. -
out.count() == 1independently guards the filter, so a change that lets test rows through fails even if the surviving row still looks right.
Output:
| check | result |
|---|---|
| normalized email |
ada@x.io ✓ |
| test row filtered | count == 1 ✓ |
Why this works — concept by concept:
- Gem-level unit test — testing one gem function in isolation makes failures point at exactly one transform, not a whole pipeline run.
- Expected-output assertion — pinning the full output (not just "no error") is what turns a smoke test into a real regression guard.
- CI enforcement — because the test is code in the repo, it runs on every pull request, so a visual edit that breaks the contract fails the build, not production.
- Independent invariants — asserting both the value and the row count catches two distinct failure modes (bad normalization vs a broken filter) separately.
- Cost — a two-row test runs in milliseconds locally, so it is cheap enough to attach to every gem that carries business rules.
Transform
Topic — data-transformation
Reusable transform and cleaning problems
Cheat sheet — Prophecy recipes
Minimal Spark pipeline (Source → Reformat → Target).
def pipeline(spark):
src = spark.read.table("sales.orders")
out = src.withColumn("total", F.col("qty") * F.col("price"))
out.write.mode("overwrite").saveAsTable("sales.orders_enriched")
Generated gem function shape (one function per gem).
def Reformat(spark, in0):
return in0.select("id", (F.col("qty") * F.col("price")).alias("total"))
Aggregate gem → groupBy/agg.
def Aggregate(spark, in0):
return in0.groupBy("customer").agg(F.sum("total").alias("revenue"))
SQL project (dbt model) gem.
select o.order_id, c.region
from {{ ref('orders') }} o
join {{ ref('customers') }} c on o.customer_id = c.customer_id
Reusable subgraph signature.
def StandardizeCustomer(spark, in0): # one input port, one output port
return in0.withColumn("email", F.lower(F.trim(F.col("email"))))
Gem unit test skeleton.
def test_gem(spark):
src = spark.createDataFrame([(1, 2, 10.0)], ["id", "qty", "price"])
out = Reformat(spark, src)
assert out.collect()[0]["total"] == 20.0
Choosing project language.
| Situation | Project type |
|---|---|
| Large distributed transforms, complex joins, ML features | Spark (PySpark/Scala) |
| Data already in the warehouse, want dbt-tested models | SQL (dbt) |
| Migrating Informatica/Alteryx mappings to open code | Spark or SQL by target engine |
| Analysts contributing reviewable pipelines | either — both emit Git-tracked code |
Frequently asked questions
What is Prophecy?
Prophecy is a low-code data engineering platform: a visual, drag-and-drop canvas for building data pipelines that compiles to open-source code. You wire together gems (transform nodes) on a graph, and Prophecy generates readable PySpark or Scala Spark for Spark pipelines, or dbt-style SQL models for warehouse pipelines, and commits that code to your own Git repository. It is a design-and-generate layer, not a runtime you send data through.
Does Prophecy lock me into its runtime?
No — that is the central design choice. The artifact Prophecy produces is standard Spark or SQL code in your Git repo; the Spark runs on any Spark cluster (Databricks, EMR, Dataproc) via spark-submit, and the SQL runs through dbt on your warehouse. Neither needs Prophecy at execution time, so if you stopped using the visual editor you would still have a working, readable, tested codebase.
What are gems in Prophecy?
Gems are the building blocks of a pipeline — each gem is one transform node on the canvas. Built-in gems cover the standard operations (Source, Target, Reformat, Aggregate, Join, Filter, OrderBy, Deduplicate, SetOperation, FlattenSchema), and each gem compiles to exactly one function in the generated code. Platform teams can also build custom gems with the Gem Builder to encode company-specific logic as reusable drag-and-drop boxes.
How does the visual-to-code round-trip work?
Every edit on the visual canvas immediately regenerates the corresponding code, and because Prophecy also parses code back into the graph, an engineer can edit the generated file in an IDE and see the canvas update. The two views cannot drift apart. Since the code is the save format and lives in Git, you get branches, commits, pull requests, reviews, and CI/CD exactly like any software repository.
Does Prophecy support SQL as well as Spark?
Yes. A Spark project compiles gems to PySpark/Scala DataFrame transformations that run on a Spark cluster, while a SQL project compiles the same style of gem canvas to dbt SQL models that run in-warehouse (Databricks SQL, Snowflake, BigQuery). The authoring experience is the same; the project language and the Fabric (execution environment) decide the dialect and engine.
Can Prophecy migrate legacy Informatica or Alteryx pipelines?
Yes — a common enterprise use case is converting legacy visual ETL (Informatica, Alteryx, Ab Initio, SSIS, DataStage) into open Spark or SQL. Prophecy's migration tooling maps the legacy mappings/workflows onto gems and generates the equivalent PySpark or dbt code in Git, so you exit the proprietary runtime while keeping a visual editor and a modern, testable, version-controlled codebase.
Practice on PipeCode
Pipecode.ai is Leetcode for Data Engineering — every Prophecy idea above, from the projection Reformat gem to the group-by rollup, the ROW_NUMBER dedup, and the unit-tested reusable subgraph, maps to a hands-on practice room where you write the Spark or SQL against real graded inputs. PipeCode pairs each reading with 450+ DE-focused problems and a real-time scoring engine, so your answer to "what code would that gem generate?" holds up under a senior interviewer's depth probes.





Top comments (0)