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 > >
