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