DEV Community

Cover image for Prophecy: Visual, Git-Backed Low-Code Spark & SQL Pipelines for the Enterprise
Gowtham Potureddi
Gowtham Potureddi

Posted on

Prophecy: Visual, Git-Backed Low-Code Spark & SQL Pipelines for the Enterprise

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.

PipeCode blog header for Prophecy — bold white headline 'Prophecy: Visual → Code' with subtitle 'Git-Backed Low-Code Spark & SQL Pipelines' and a stylised visual-canvas-compiles-to-code scene on a dark gradient with purple, green, orange, and blue accents and a small pipecode.ai attribution.

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


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)
Enter fullscreen mode Exit fullscreen mode

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 (Spark functions calls 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 DataFrame at 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.

Iconographic Prophecy canvas diagram — source, reformat, aggregate, join and target gems wired into a pipeline DAG, each gem compiling to a function in the generated code, with a per-gem data-preview chip.

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")
Enter fullscreen mode Exit fullscreen mode

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"),
    )
Enter fullscreen mode Exit fullscreen mode

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
  1. 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.
  2. total multiplies two columns; year extracts the year from a timestamp; email is normalized by nesting lower around trim.
  3. Because everything is one select, Spark's Catalyst optimizer fuses the expressions into a single project stage — no extra shuffle or scan.
  4. Naming each derived column with .alias(...) is what makes the gem's output schema explicit and lets the next gem reference total and year by name.

Output:

order_id qty price total year email
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 chained withColumn calls that stack projections.
  • Column expressions — writing qty * price in the gem and letting Prophecy emit F.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

Practice →

ETL Topic — etl ETL pipeline-design problems

Practice →


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.py in 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."

Iconographic Prophecy round-trip diagram — a visual canvas and a code file kept in sync by a two-way arrow, both committed to a Git branch with commit, PR and merge glyphs, and a note that the generated code runs without Prophecy.

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"),
           )
    )
Enter fullscreen mode Exit fullscreen mode

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)
    )
Enter fullscreen mode Exit fullscreen mode

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
  1. groupBy("customer") collapses each customer's rows into one group, triggering one shuffle keyed on customer.
  2. sum("total") accumulates revenue and countDistinct(to_date(order_ts)) counts unique calendar days, both computed in the same aggregation pass.
  3. The post-aggregation filter(revenue > 100) is a HAVING clause: it runs on the grouped rows, dropping linus while keeping ada and grace.
  4. 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 customer so all of one customer's rows meet on one executor; this is the single shuffle the whole transform costs.
  • Multiple aggregates, one pass — sum and countDistinct are computed together, so two metrics cost one scan and one shuffle, not two.
  • HAVING via post-filter — filtering after agg is the Spark equivalent of SQL HAVING; the predicate references the aggregated revenue, which is only defined after the group.
  • countDistinct semantics — counting distinct to_date(order_ts) answers "how many active days," a different question from count(*); 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

Practice →

Transform Topic — data-transformation Group-by and HAVING transform problems

Practice →


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 DataFrame transformations; 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.partitions you would tune by hand.

Iconographic Prophecy execution diagram — a Fabric config card connecting to a Databricks Spark cluster and a SQL warehouse, with a Spark project compiling to PySpark and a SQL project compiling to dbt models, deployed as scheduled jobs.

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")
Enter fullscreen mode Exit fullscreen mode
-- 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
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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
  1. partition by order_id groups the duplicate rows for each order together so ranking is per-key.
  2. order by updated_at desc puts the freshest row first; the tie-break ingest_seq desc makes the winner deterministic when two rows share the same updated_at.
  3. row_number() assigns 1 to the winner of each partition, 2..N to the rest.
  4. where rn = 1 keeps exactly one row per order_id, and except (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_id turns "dedupe per key" into a ranking problem solved inside each partition, no self-join required.
  • Deterministic tie-break — adding ingest_seq desc to the order by removes the nondeterminism of ties, so re-running yields the same survivor every time.
  • row_number vs rank — row_number guarantees a single winner (no ties at rank 1), which is exactly the "keep one" requirement; rank could 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_id plus 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

Practice →

Optimization Topic — optimization Shuffle and partition tuning problems

Practice →


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 revenue back 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.

Iconographic Prophecy reuse diagram — a group of gems packaged as a reusable subgraph component, a gem unit test with sample input and expected output, and a column-level lineage graph tracing a field across pipelines.

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"))
    )
Enter fullscreen mode Exit fullscreen mode

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 email
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
Enter fullscreen mode Exit fullscreen mode

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
  1. The test builds a tiny DataFrame with one good row and one @test.local row, then calls the exact gem function under test.
  2. StandardizeCustomer trims and lowercases the email, uppercases the name, and filters the @test.local row out.
  3. The assertion pins the whole expected output — one row, {"ADA": "ada@x.io"} — so any regression (a missing lower, a broken filter) fails the assertion.
  4. out.count() == 1 independently 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

Practice →

Quality Topic — data-quality Unit-test and data-validation problems

Practice →


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")
Enter fullscreen mode Exit fullscreen mode

Generated gem function shape (one function per gem).

def Reformat(spark, in0):
    return in0.select("id", (F.col("qty") * F.col("price")).alias("total"))
Enter fullscreen mode Exit fullscreen mode

Aggregate gem → groupBy/agg.

def Aggregate(spark, in0):
    return in0.groupBy("customer").agg(F.sum("total").alias("revenue"))
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

Reusable subgraph signature.

def StandardizeCustomer(spark, in0):   # one input port, one output port
    return in0.withColumn("email", F.lower(F.trim(F.col("email"))))
Enter fullscreen mode Exit fullscreen mode

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
Enter fullscreen mode Exit fullscreen mode

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.

Practice Spark SQL problems now →
ETL pipeline drills →

Top comments (0)