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]