Hi Eric,

I think you found a real bug, and also a gap in the test coverage.
I am curious what your fix looks like. So yes, please open a Jira and send
a PR, that would be really helpful.

Thanks -Andreas


On Sun, Sep 13, 2026 at 7:37 PM Eric Smith <[email protected]> wrote:

> Hi all,
>
> I've been poking at a hobby Scala 3 client for Declarative Pipelines over
> Spark Connect (mostly as a way to learn the SDP internals properly), which
> means I spend a lot of time in corners of the pipelines code most people
> probably don't visit. I ran into something last week that I can't explain
> away, and before I file anything I'd really like a sanity check from folks
> who actually know this code. There is a non-zero chance I'm
> misunderstanding how re-registration is supposed to work.
>
> What I see
>
> A two-dataset pipeline, stock apache/spark:4.2.0, official spark-pipelines
> CLI:
>
> ```python
> @dp.materialized_view
> def a():
>     return spark.sql("SELECT * FROM VALUES (1,'north'),(2,'south') AS
> t(id, name)")
>
>
> @dp.materialized_view
> def b():
>     return spark.read.table("a").withColumn("m", F.upper(F.col("name")))
>
>
> @dp.materialized_view
> def c():   # control
>     return spark.read.table("a").select("id", "name",
> F.upper(F.col("name")).alias("m"))
> ```
>
> Run 1 is fine: a completes, then b and c run, both get 2 rows. Run the
> same command again and b comes back empty. The server's own ordering log
> shows why. On run 2, b starts at the same instant as a instead of after it,
> read the table DatasetManager just truncated, and writes 0 rows. c still
> waits for a and is fine on both runs. The run reports COMPLETED, exit 0, no
> warning, and dry-run doesn't catch it either.
>
> So the difference between MVs c and b is select vs withColumn
> respectively, and issues only occure if the upstream table already exists
> in the catalog.
>
> What I think is happening
>
> SparkConnectPlanner.transformWithColumns does eager analysis of its child
> (Dataset.ofRows(session, child). On a clean catalog that throws and the
> fallback keeps the UnresolvedRelation. But once a exists, the eager
> analysis succeeds at DefineFlow time, the read arrives at FlowAnalysis
> already resolved, and FlowAnalysis.analyze only registers dependencies for
> case u: UnresolvedRelation. Now, no edge gets recorded. A plain
> df.explain(extended=True) from a Connect session shows it: with a present,
> the withColumn variant's Parsed plan is already Relation
> spark_catalog.default.a, the select variant's is still 'UnresolvedRelation
> [a].
>
> I could be wrong about the mechanism; the observation is what I'm
> confident about.
>
> The failing test
>
> I lifted test("referencing internal datasets") from PythonPipelineSuite
> (line 351 on v4.2.0) and changed three things: the setupSql hook
> pre-creates src so the graph registers against an existing table (this is
> my second run state from above), a gets a .withColumn, and I kept a .select
> sibling as the control.
>
> test("referencing internal datasets that already exist in the catalog") {
>   val graph = buildGraph(
>     """
>       |from pyspark.sql import functions as F
>       |
>       |@dp.materialized_view
>       |def src():
>       |  return spark.range(5)
>       |
>       |@dp.materialized_view
>       |def a():
>       |  return spark.read.table("src").withColumn("y", F.col("id") + 1)
>       |
>       |@dp.materialized_view
>       |def c():
>       |  return spark.read.table("src").select("id")
>       |""".stripMargin,
>     // the precondition: `src` already exists, as it does on every run
> after the first
>     setupSql = Some("CREATE TABLE spark_catalog.default.src AS SELECT *
> FROM range(5)")
>   ).resolve(sessionCaseSensitive).validate(sessionCaseSensitive)
>
>
>   assert(graph.resolvedFlow(graphIdentifier("c")).inputs ==
> Set(graphIdentifier("src")))
>   assert(graph.resolvedFlow(graphIdentifier("a")).inputs ==
> Set(graphIdentifier("src")))
> }
>
> I ran that on master with the test dropped into PythonPipelineSuite right
> after test("referencing internal datasets"), via build/sbt
> 'connect/testOnly *PythonPipelineSuite -- -z "already exist"'. The c
> assertion passes and the a one fails:
>
> - referencing internal datasets that already exist in the catalog ***
> FAILED *** (4 seconds, 198 milliseconds)
>   Set() did not equal Set(`spark_catalog`.`default`.`src`)
> (PythonPipelineSuite.scala:429)
>   Analysis:
>   Set$EmptySet$(missingInLeft: [`spark_catalog`.`default`.`src`])
>
> The test it was lifted from still passes on that same tree, so this isn't
> the harness.
>
> As far as I can tell none of the existing tests register a flow while its
> upstream tables already exists (every inputs assertion I found runs against
> a clean catalog).
>
> So, two questions really: is this expectation wrong (maybe re-registering
> over already-materialized tables isn't a supported flow and I'm holding it
> wrong)? And if it's real, how would you want it fixed? I have a small patch
> that plans flow relations without the eager analysis, but is it the right
> path? Happy to file a JIRA and open a PR with the test if that's useful.
>
> Sorry for the length, and thanks for reading. Even a "this is known /
> here's the doc you missed" would save me from chasing my tail.
>
> Eric
>
>

Reply via email to