[
https://issues.apache.org/jira/browse/SPARK-59395?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Paul Sedra updated SPARK-59395:
-------------------------------
Description:
h1. Summary
Spark Declarative Pipelines (SDP) does not expose a deterministic way to
associate an individual SDP flow execution with the Spark SQL executions and
jobs it produces, or with their stages and tasks. Pipeline-level attribution
can be achieved using existing Spark execution metadata, but that context does
not distinguish individual flows.
SDP knows the active flow internally, but that identity is not exposed through
the public observability boundary. External tools therefore cannot
deterministically associate an SDP flow execution with the Spark work it
produces.
h2. Current behavior
A caller can attach an opaque pipeline-run identifier through Spark Connect
session/operation metadata or job tags. This establishes Spark work →
containing pipeline run, but not Spark work → the SDP flow and execution
attempt that created it. The public SDP Spark Connect StartRun request has no
flow-execution identity, and an enclosing ExecutePlanRequest tag is not
guaranteed to propagate through asynchronous per-flow execution. Pipeline
events provide status text and timestamps, not structured flow-to-execution
links. SQL text, query plans, and timestamp matching are not authoritative.
h2. Minimal example
For two flows in one run: pipeline_run = R1; silver_orders → SQL E1 → Spark job
J1; gold_orders → SQL E2 → Spark job J2. Current tagging can show J1 → R1 and
J2 → R1, but cannot deterministically show which flow produced J1 or J2.
h2. Runtime evidence
The OSS runtime knows the active flow at the relevant boundary:
GraphExecution.planAndStartFlow(flow) → FlowExecution.executeAsync → batch or
streaming execution. See
[GraphExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala#L79|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala#L957-L980]]
and
[FlowExecution.scala|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala#L141].
That internal context is not exposed as a supported public SDP callback,
protocol field, or Spark
execution-metadata contract.
h2. Expected behavior
Expose enough stable semantic identity for an external observer to determine:
- which logical SDP flow is executing;
- which individual execution or attempt is being observed; and
- which Spark SQL executions and/or Spark jobs belong to that execution,
allowing existing Spark
execution relationships to provide stage/task attribution.
The behavior should cover SDP batch and streaming flows. The implementation and
API shape are intentionally left to Spark maintainers; protocol metadata,
execution tags, structured events, or listener/event-log metadata are possible
mechanisms, not requirements.
h2. Acceptance criteria
# An SDP flow has an externally observable logical identity and an identity
for an individual execution or attempt.
# External tools can deterministically correlate that execution with its Spark
SQL executions and/or jobs, allowing existing Spark execution relationships to
provide stage/task attribution, without parsing SQL, plans, logs, or timestamps.
# Attribution remains correct for multiple flows and attempts, including batch
and streaming flows, while existing clients that do not use the new metadata
remain compatible.
h2. References
- SPARK-51727: SPIP: Declarative Pipelines
- SPARK-44591: Add jobTags to SparkListenerSQLExecutionStart
- SPARK-44612: Use jobTags in SparkListenerSQLExecutionStart to get SQL
Execution ID for Spark UI Connect page
-
[GraphExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala]|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala],
[FlowExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/FlowExecution.scala]|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/FlowExecution.scala],
and
[pipelines.proto|[https://github.com/apache/spark/blob/master/sql/connect/common/src/main/protobuf/spark/connect/pipelines.proto]|https://github.com/apache/spark/blob/master/sql/connect/common/src/main/protobuf/spark/connect/pipelines.proto]
was:
h1. Summary
Spark Declarative Pipelines (SDP) does not expose a deterministic way to
associate an individual SDP flow execution with the Spark SQL executions and
jobs it produces, or with their stages and tasks. Pipeline-level attribution
can be achieved using existing Spark execution metadata, but that context does
not distinguish individual flows.
SDP knows the active flow internally, but that identity is not exposed through
the public observability boundary. External tools therefore cannot
deterministically associate an SDP flow execution with the Spark work it
produces.
h2. Current behavior
A caller can attach an opaque pipeline-run identifier through Spark Connect
session/operation metadata or job tags. This establishes Spark work →
containing pipeline run, but not Spark work → the SDP flow and execution
attempt that created it. The public SDP Spark Connect StartRun request has no
flow-execution identity, and an enclosing ExecutePlanRequest tag is not
guaranteed to propagate through asynchronous per-flow execution. Pipeline
events provide status text and timestamps, not structured flow-to-execution
links. SQL text, query plans, and timestamp matching are not authoritative.
h2. Minimal example
For two flows in one run: pipeline_run = R1; silver_orders → SQL E1 → Spark job
J1; gold_orders → SQL E2 → Spark job J2. Current tagging can show J1 → R1 and
J2 → R1, but cannot deterministically show which flow produced J1 or J2.
h2. Runtime evidence
The OSS runtime knows the active flow at the relevant boundary:
GraphExecution.planAndStartFlow(flow) → FlowExecution.executeAsync → batch or
streaming execution. See
[GraphExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala#L79|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala#L957-L980]
[]|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala#L957-L980]
and
[FlowExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/FlowExecution.scala#L141]|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/FlowExecution.scala#L1225-L1287].
That internal context is not exposed as a supported public SDP callback,
protocol field, or Spark
execution-metadata contract.
h2. Expected behavior
Expose enough stable semantic identity for an external observer to determine:
- which logical SDP flow is executing;
- which individual execution or attempt is being observed; and
- which Spark SQL executions and/or Spark jobs belong to that execution,
allowing existing Spark
execution relationships to provide stage/task attribution.
The behavior should cover SDP batch and streaming flows. The implementation and
API shape are intentionally left to Spark maintainers; protocol metadata,
execution tags, structured events, or listener/event-log metadata are possible
mechanisms, not requirements.
h2. Acceptance criteria
# An SDP flow has an externally observable logical identity and an identity
for an individual execution or attempt.
# External tools can deterministically correlate that execution with its Spark
SQL executions and/or jobs, allowing existing Spark execution relationships to
provide stage/task attribution, without parsing SQL, plans, logs, or timestamps.
# Attribution remains correct for multiple flows and attempts, including batch
and streaming flows, while existing clients that do not use the new metadata
remain compatible.
h2. References
- [SPARK-51727: SPIP: Declarative
Pipelines|https://issues.apache.org/jira/browse/SPARK-51727]
- [SPARK-44591: Add jobTags to
SparkListenerSQLExecutionStart|https://issues.apache.org/jira/browse/SPARK-44591]
- [SPARK-44612: Use jobTags in SparkListenerSQLExecutionStart to get SQL
Execution ID for Spark UI Connect
page|https://issues.apache.org/jira/browse/SPARK-44612]
-
[GraphExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala]|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala],
[FlowExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/FlowExecution.scala]|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/FlowExecution.scala],
and
[pipelines.proto|[https://github.com/apache/spark/blob/master/sql/connect/common/src/main/protobuf/spark/connect/pipelines.proto]|https://github.com/apache/spark/blob/master/sql/connect/common/src/main/protobuf/spark/connect/pipelines.proto]
> [SDP] Expose deterministic flow-to-Spark execution attribution
> --------------------------------------------------------------
>
> Key: SPARK-59395
> URL: https://issues.apache.org/jira/browse/SPARK-59395
> Project: Spark
> Issue Type: Improvement
> Components: Declarative Pipelines
> Affects Versions: 4.1.3
> Reporter: Paul Sedra
> Priority: Minor
>
> h1. Summary
> Spark Declarative Pipelines (SDP) does not expose a deterministic way to
> associate an individual SDP flow execution with the Spark SQL executions and
> jobs it produces, or with their stages and tasks. Pipeline-level attribution
> can be achieved using existing Spark execution metadata, but that context
> does not distinguish individual flows.
> SDP knows the active flow internally, but that identity is not exposed
> through the public observability boundary. External tools therefore cannot
> deterministically associate an SDP flow execution with the Spark work it
> produces.
> h2. Current behavior
> A caller can attach an opaque pipeline-run identifier through Spark Connect
> session/operation metadata or job tags. This establishes Spark work →
> containing pipeline run, but not Spark work → the SDP flow and execution
> attempt that created it. The public SDP Spark Connect StartRun request has no
> flow-execution identity, and an enclosing ExecutePlanRequest tag is not
> guaranteed to propagate through asynchronous per-flow execution. Pipeline
> events provide status text and timestamps, not structured flow-to-execution
> links. SQL text, query plans, and timestamp matching are not authoritative.
> h2. Minimal example
> For two flows in one run: pipeline_run = R1; silver_orders → SQL E1 → Spark
> job J1; gold_orders → SQL E2 → Spark job J2. Current tagging can show J1 → R1
> and J2 → R1, but cannot deterministically show which flow produced J1 or J2.
> h2. Runtime evidence
> The OSS runtime knows the active flow at the relevant boundary:
> GraphExecution.planAndStartFlow(flow) → FlowExecution.executeAsync → batch or
> streaming execution. See
> [GraphExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala#L79|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala#L957-L980]]
> and
> [FlowExecution.scala|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala#L141].
> That internal context is not exposed as a supported public SDP callback,
> protocol field, or Spark
> execution-metadata contract.
> h2. Expected behavior
> Expose enough stable semantic identity for an external observer to determine:
> - which logical SDP flow is executing;
> - which individual execution or attempt is being observed; and
> - which Spark SQL executions and/or Spark jobs belong to that execution,
> allowing existing Spark
> execution relationships to provide stage/task attribution.
> The behavior should cover SDP batch and streaming flows. The implementation
> and API shape are intentionally left to Spark maintainers; protocol metadata,
> execution tags, structured events, or listener/event-log metadata are
> possible mechanisms, not requirements.
> h2. Acceptance criteria
> # An SDP flow has an externally observable logical identity and an identity
> for an individual execution or attempt.
> # External tools can deterministically correlate that execution with its
> Spark SQL executions and/or jobs, allowing existing Spark execution
> relationships to provide stage/task attribution, without parsing SQL, plans,
> logs, or timestamps.
> # Attribution remains correct for multiple flows and attempts, including
> batch and streaming flows, while existing clients that do not use the new
> metadata remain compatible.
> h2. References
> - SPARK-51727: SPIP: Declarative Pipelines
> - SPARK-44591: Add jobTags to SparkListenerSQLExecutionStart
> - SPARK-44612: Use jobTags in SparkListenerSQLExecutionStart to get SQL
> Execution ID for Spark UI Connect page
> -
> [GraphExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala]|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/GraphExecution.scala],
>
> [FlowExecution.scala|[https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/FlowExecution.scala]|https://github.com/apache/spark/blob/master/sql/pipelines/src/main/scala/org/apache/spark/sql/pipelines/graph/FlowExecution.scala],
> and
> [pipelines.proto|[https://github.com/apache/spark/blob/master/sql/connect/common/src/main/protobuf/spark/connect/pipelines.proto]|https://github.com/apache/spark/blob/master/sql/connect/common/src/main/protobuf/spark/connect/pipelines.proto]
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]