jerrypeng commented on code in PR #57092:
URL: https://github.com/apache/spark/pull/57092#discussion_r3590947981


##########
PIPELINED_SHUFFLE_DEPENDENCY_SPEC.md:
##########
@@ -0,0 +1,409 @@
+# Pipelined Shuffle Dependency & Concurrent Stage Scheduling
+
+A spec for running data-dependent stages of a single job concurrently, 
connected by a shuffle the
+consumer reads incrementally.
+
+---
+
+## 1. Motivation
+
+Today, when two stages of a job are connected by a shuffle, they run one after 
another: the consumer
+stage does not start until the producer's shuffle output is fully 
materialized. (Independent stages —
+ones not connected by a shuffle — already run concurrently.) Some workloads 
need even
+shuffle-connected stages to run **concurrently**, with the consumer reading 
the producer's output
+**as it is produced** rather than after the producer finishes. This spec 
introduces the scheduler
+primitives to express and run that.
+
+"Run these stages concurrently" and "the connecting shuffle is incremental" 
are the same decision
+seen from two sides: co-scheduling a producer and consumer is only useful if 
the edge is readable
+before the producer completes.
+
+---
+
+## 2. Primitives
+
+### 2.1 Pipelined shuffle dependency (PSD)
+
+A shuffle dependency declared **incrementally readable**: a consumer stage may 
begin reading its
+output while the producer stage is still running.
+
+- It is a shuffle dependency (has a `shuffleId`, partitioner, map/reduce 
sides); the *pipelined*
+  property is a binding part of the scheduler contract, not an advisory hint. 
("Pipelined",
+  "incremental", and "transient" all describe this same edge in this doc — 
read as its output being
+  streamed to the consumer as produced, not stored.)
+- The property is set during **physical planning** (an execution concern, not 
a logical-plan one)
+  and carried into the `ShuffleDependency` the `DAGScheduler` reads at 
stage-creation time.
+- It is also the **per-dependency selector** for the shuffle implementation: 
the shuffle layer maps a
+  pipelined dependency to an incremental `ShuffleManager` and everything else 
to the default, so one
+  job with both regular and pipelined groups uses the right implementation for 
each. The scheduler
+  construct stays generic; the shuffle implementation stays pluggable.
+  - For example, an additional conf (e.g. `spark.shuffle.manager.incremental`) 
can be introduced to
+    specify the incremental shuffle implementation used for pipelined shuffle 
dependencies, alongside
+    the existing `spark.shuffle.manager` for the default.
+
+**What the subtype inherits vs. what the scheduler skips.** A pipelined 
shuffle dependency is modeled
+as a `ShuffleDependency` subtype, which cleanly divides `ShuffleDependency`'s 
behavior in two — an
+invariant every component that handles the type must preserve:
+- *Inherited — it is treated as a real shuffle dependency:* (1) 
**stage-boundary splitting** —
+  `getShuffleDependenciesAndResourceProfiles` matches `case shuffleDep: 
ShuffleDependency`
+  (`DAGScheduler.scala:882`), so the subtype splits producer and consumer into 
separate stages
+  automatically; this is what §3 ("introduces a shuffle boundary exactly as a 
regular one does")
+  relies on. (2) **transport** — `shuffleId`, `ShuffleWriter`/`ShuffleReader`, 
and `ShuffleManager`
+  registration, all from the base constructor.
+- *Scheduler-skipped — it must not be treated as a materialized shuffle for 
the post-boundary
+  lifecycle:* `MapOutputTracker` registration (§6), `shuffleIdToMapStage` 
reuse (§4), and
+  executor-loss output-removal-then-resubmit (§6). This is clean rather than 
"inherit then disable":
+  all three are `DAGScheduler`-driven at `createShuffleMapStage` 
(`DAGScheduler.scala:667-673`) and
+  gated per dependency, so a pipelined dependency is simply never opted into 
them.
+- *Match-site invariant:* a marker subtype relies on every existing `case _: 
ShuffleDependency` site
+  meaning "materialized" for the new type. The reuse/tracking/loss sites are 
safe because they are
+  scheduler-gated as above; the remaining sites were audited and treat a 
pipelined dependency as an
+  ordinary materialized shuffle wherever the subtype does not specifically 
diverge. Preserving this
+  is an invariant for implementers, not something to rediscover.
+
+A **regular shuffle dependency (RSD)** is an ordinary shuffle dependency: its 
output must be fully
+materialized before any consumer reads it.
+
+Note: the name *pipelined* is deliberately chosen over *streaming*. The 
property is that a consumer
+reads producer output as it is produced — software-pipelining of dependent 
stages — a general
+execution capability. Streaming / real-time mode (RTM) is the first caller, 
but nothing about the
+primitive is streaming-specific.
+
+### 2.2 Pipelined group (PG)
+
+The set of stages connected to one another through pipelined edges — the 
connected component of the
+stage DAG when only pipelined edges are considered.
+
+- A stage with no incident pipelined edge is a **singleton group** and behaves 
exactly as a normal
+  stage today.
+- The group — not the edge or the individual stage — is the unit of 
**admission**, **slot
+  checking**, **completion**, and **failure**.
+
+**External input of PG:** a regular shuffle dependency whose consumer is in 
the group and whose
+producer is not — i.e. a normal materialized parent of the group.
+
+**Outputs of PG:** a `ResultStage` **may** be a member of a pipelined group — 
a pipelined consumer
+that produces the job's result (its commit handling is §5) — but a PG **need 
not** contain one.
+When it does not, its outputs to the rest of the DAG are regular 
(materialized) shuffle edges from
+its members to consumers outside the group, and there may be **more than one** 
— e.g. two members
+each feeding a distinct downstream stage. Such a PG is an interior component 
whose materialized
+outputs feed downstream groups (§7); in that respect it behaves like a barrier 
stage embedded in a
+larger DAG — an interior unit with regular-shuffle edges in and out.
+
+**Example.** Take a job with stage DAG `Z -> A -> B -> C`, where `Z -> A` and 
`B -> C` are regular
+edges and `A -> B` is pipelined. Grouping over pipelined edges gives one 
non-singleton group
+`PG = {A, B}` (`Z` and `C` are singletons). `Z`'s output is the group's 
external input, which the
+group waits for; `A` and `B` run concurrently, `B` reading `A`'s output as `A` 
produces it; and
+`B`'s output to `C` is a materialized output of the group that `C` consumes by 
normal
+materialize-before-read sequencing.
+
+---
+
+## 3. Group formation
+
+- **Stage decomposition is unchanged.** A pipelined dependency introduces a 
shuffle boundary exactly
+  as a regular one does; the set of stages and their partitioning are 
identical. The pipelined
+  property changes only *when* stages run relative to one another, never *how 
the plan is cut into
+  stages*.
+- **Group = connected component over pipelined edges.** As stages are created, 
two stages joined by a
+  pipelined edge are placed in the same group; the group is the transitive 
closure.
+- **Every stage belongs to exactly one group** (singletons included). Group 
membership is fixed at
+  stage-creation time.
+- **A DAG may contain multiple pipelined groups.** Disjoint sets of pipelined 
edges form distinct
+  groups, and a group's materialized output may feed another group (§7). 
Independent groups (no
+  dependency path between them) may be admitted and run concurrently, subject 
to the cross-group
+  slot arbitration in §4.1; a group that consumes another group's materialized 
output waits for it,
+  by the normal materialize-before-read sequencing on that regular edge. This 
mirrors how a DAG may
+  contain multiple barrier stages.
+
+---
+
+## 4. Scheduling & admission
+
+- **A pipelined edge is non-sequencing.** The consumer of a pipelined 
dependency does not wait for
+  its producer to materialize. (A regular-dependency consumer still waits — 
the default behavior.)
+- **Group readiness.** A group is ready to be admitted when every external 
input (its regular
+  materialized parents) is available — the same precondition a normal stage 
has today, lifted to the
+  group. Pipelined parents inside the group impose no readiness precondition.
+- **Gang admission (all-or-nothing).** A pipelined group is admitted only if 
the cluster can currently
+  run all tasks of all member stages **concurrently**; admission then submits 
every member stage at
+  once. A pipelined group is never left with some members running while others 
wait on slots the
+  running members occupy.
+  - *A non-pipelined (singleton) group is unaffected:* it is admitted exactly 
as a normal stage is
+    today — submitted when its parents are available, filling slots 
incrementally, with no all-at-once
+    requirement. Gang admission applies only to a group of two or more stages 
connected by pipelined
+    edges.
+- **Slot check.** The group's aggregate concurrent-task demand — the sum of 
`numTasks` over member
+  stages — is compared against the number of **currently-free** slots: the 
cluster's total
+  concurrent-task capacity minus the tasks already running. If the group needs 
more than are free,
+  admission fails fast, since the group could never become fully co-resident. 
(Free rather than
+  total is essential once more than one group can be admitted; §4.1 explains 
why.)
+  - *How capacity is measured:* the total concurrent-task capacity is the 
value barrier's slot check
+    uses (`sc.maxNumConcurrentTasks` — for the group's resource profile, the 
number of task slots
+    across active executors, i.e. cores divided by cores-per-task), not 
`spark.default.parallelism`
+    (a default partition count, unrelated to how many tasks can run 
concurrently).
+- **Single resource profile per group (v1).** A resource profile is the 
executor/task resource
+  requirement (cores, memory, GPUs, ...) a stage runs under; the number of 
concurrent slots is
+  defined *per profile* (a cluster may run many concurrent CPU tasks but few 
GPU tasks). Comparing
+  one demand against one capacity is therefore only well-defined when all 
member stages of a group
+  share a single profile. v1 requires that and rejects a mixed-profile group 
(fail-fast, §9);
+  per-profile accounting — checking each profile's demand against that 
profile's own capacity — is a
+  follow-up, not needed for the streaming shapes whose members share the 
default profile.
+- **Co-residency.** Once admitted, all member stages of a group are 
simultaneously running.
+- **Single ownership and 1:1 (v1).** In v1 a pipelined producer feeds exactly 
one consumer and
+  belongs to exactly one job. These are two separate restrictions. They follow 
from what a pipelined
+  shuffle *is* (a transient edge), not from how any particular caller builds 
its DAG — the pipelined
+  dependency and PG are generic `DAGScheduler` constructs that code can depend 
on directly, not only
+  through SQL / real-time mode:
+  - *No cross-job / cross-time reuse — intrinsic.* A pipelined shuffle is 
transient: a once-through
+    live stream with no retained, addressable output. So, unlike a regular 
shuffle, there is nothing
+    for a second job to reuse — reuse across jobs is unsound for it.
+    - This must be enforced, not assumed: from the scheduler's perspective a 
shuffle-map stage *can*

Review Comment:
   let me clarify this.  The behavior is a second job over the same pipelined 
dependency is an error, rejected fail-fast — never silent reuse, never silent 
recompute.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to