Eric Smith created SPARK-59663:
----------------------------------
Summary: SDP: `withColumn` on an in-graph table read silently
drops the dependency edge once the upstream table exists
Key: SPARK-59663
URL: https://issues.apache.org/jira/browse/SPARK-59663
Project: Spark
Issue Type: Bug
Components: Declarative Pipelines
Affects Versions: 4.1.3, 4.1.2, 4.1.1, 4.2.0, 4.1.0
Environment: Image: docker.io/apache/spark:4.2.0
Spark Version: 4.2.0
Client: /opt/spark/bin/spark-pipelines
Server: start-connect-server.sh --conf spark.sql.warehouse.dir=/tmp/wh --conf
spark.sql.catalogImplementation=in-memory
Format: default parquet
reproduction
`spark-pipeline.yaml`:
```yaml
name: edge_drop_repro
storage: file:///tmp/sdp-jira-abc-storage
catalog: spark_catalog
database: repro
libraries:
- glob:
include: transformations/**
```
`transformations/pipeline.py`:
```python
from pyspark import pipelines as dp
from pyspark.sql import functions as F
@dp.materialized_view
def a():
return spark.sql("SELECT * FROM VALUES (1, 'north'), (2, 'south') AS t(id,
name)")
@dp.materialized_view
def b():
# withColumn -> proto WithColumns -> SparkConnectPlanner.transformWithColumns
return spark.read.table("a").withColumn("m", F.upper(F.col("name")))
@dp.materialized_view
def c():
# CONTROL: select -> proto Project -> transformProject, no eager analysis
return spark.read.table("a").select("id", "name",
F.upper(F.col("name")).alias("m"))
```
```bash
# the database must pre-exist: spark.sql("CREATE DATABASE IF NOT EXISTS
spark_catalog.repro")
/opt/spark/bin/spark-pipelines run --remote sc://localhost:15002 --spec
spark-pipeline.yaml # run 1
/opt/spark/bin/spark-pipelines run --remote sc://localhost:15002 --spec
spark-pipeline.yaml # run 2
```
**Expected:** run 2 identical to run 1 — `b` waits for `a`, then holds 2 rows.
**Run 1 (clean catalog) — correct.** `TriggeredGraphExecution` ordering:
```
26/09/19 19:28:12 INFO TriggeredGraphExecution: Starting flow
`spark_catalog`.`repro`.`a`
26/09/19 19:28:14 INFO TriggeredGraphExecution: Flow
`spark_catalog`.`repro`.`a` has COMPLETED in TriggeredFlowExecution.
26/09/19 19:28:14 INFO TriggeredGraphExecution: Starting flow
`spark_catalog`.`repro`.`c`
26/09/19 19:28:14 INFO TriggeredGraphExecution: Starting flow
`spark_catalog`.`repro`.`b`
26/09/19 19:28:15 INFO TriggeredGraphExecution: Flow
`spark_catalog`.`repro`.`c` has COMPLETED in TriggeredFlowExecution.
26/09/19 19:28:15 INFO TriggeredGraphExecution: Flow
`spark_catalog`.`repro`.`b` has COMPLETED in TriggeredFlowExecution.
```
Rows: `a`=2, `b`=2, `c`=2.
```
ROWS a count=2
ROW a Row(id=1, name='north')
ROW a Row(id=2, name='south')
ROWS b count=2
ROW b Row(id=1, name='north', m='NORTH')
ROW b Row(id=2, name='south', m='SOUTH')
ROWS c count=2
ROW c Row(id=1, name='north', m='NORTH')
ROW c Row(id=2, name='south', m='SOUTH')
```
**Run 2 (identical command) — `b` loses its edge:**
```
26/09/19 19:28:33 INFO TriggeredGraphExecution: Starting flow
`spark_catalog`.`repro`.`b`
26/09/19 19:28:33 INFO TriggeredGraphExecution: Starting flow
`spark_catalog`.`repro`.`a`
26/09/19 19:28:34 INFO TriggeredGraphExecution: Flow
`spark_catalog`.`repro`.`b` has COMPLETED in TriggeredFlowExecution.
26/09/19 19:28:34 INFO TriggeredGraphExecution: Flow
`spark_catalog`.`repro`.`a` has COMPLETED in TriggeredFlowExecution.
26/09/19 19:28:34 INFO TriggeredGraphExecution: Starting flow
`spark_catalog`.`repro`.`c`
26/09/19 19:28:35 INFO TriggeredGraphExecution: Flow
`spark_catalog`.`repro`.`c` has COMPLETED in TriggeredFlowExecution.
```
Rows: `a`=2, **`b`=0**, `c`=2.
```
ROWS a count=2
ROW a Row(id=1, name='north')
ROW a Row(id=2, name='south')
ROWS b count=0
ROWS c count=2
ROW c Row(id=1, name='north', m='NORTH')
ROW c Row(id=2, name='south', m='SOUTH')
```
Reporter: Eric Smith
`SparkConnectPlanner.transformWithColumns` analyzes its child eagerly, and
`PipelinesHandler.defineFlow` runs a flow's relation through the planner at
`DefineFlow` time.
Once the upstream table exists in the catalog, the eager analysis succeeds,
consumes the `UnresolvedRelation`, and `FlowAnalysis` (which discovers a flow's
dependencies only by matching UnresolvedRelation) records no dependency. The
downstream flow is then scheduled concurrently with its own upstream
The first run of a pipeline is always correct. The eager analysis fails on a
clean catalog and the `catch` preserves the read. So ths passes every smoke
test and fails on the second scheduled run.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]