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


##########
sql/core/src/main/scala/org/apache/spark/sql/execution/streaming/runtime/IncrementalExecution.scala:
##########
@@ -663,7 +664,91 @@ class IncrementalExecution(
     }
   }
 
-  override def preparations: Seq[Rule[SparkPlan]] = state +: super.preparations
+  /**
+   * For a Real-Time Mode batch, mark the shuffle exchanges as pipelined so 
the DAGScheduler
+   * co-schedules a stateful query's producer (source scan) and consumer 
(stateful operator) stages
+   * as one pipelined group -- records stream through a transient shuffle 
instead of the consumer
+   * waiting for the producer to fully materialize. The exchange carries the 
decision as a field
+   * (see ShuffleExchangeExec.pipelined); the PipelinedShuffleDependency it 
then builds is the whole
+   * opt-in -- routing to the streaming shuffle manager and pipelined-group 
scheduling both follow
+   * from that dependency type.
+   *
+   * Real-Time Mode is detected structurally by a RealTimeStreamScanExec leaf 
(there is no
+   * RTM-specific plan flag). Inert for a non-RTM batch, so the ordinary 
microbatch path is
+   * unchanged.
+   *
+   * Marks every eligible shuffle exchange on the streaming path, so a plan 
with several pipelined
+   * shuffles in a chain (e.g. two repartitions, or a repartition feeding a 
keyed stateful operator)
+   * is handled: each becomes a PipelinedShuffleDependency and the whole 
all-pipelined job is
+   * co-scheduled as one pipelined group (the DAGScheduler treats an 
all-pipelined job's stage graph
+   * as a single group). There is no shuffle-count restriction. An exchange 
whose subtree does not
+   * reach the real-time scan is skipped -- the static side of a broadcast 
stream-static join must
+   * materialize, because it runs to completion rather than streaming. A 
partitioning the pipelined
+   * path cannot serve, such as range partitioning, is rejected up front by 
RealTimeModeAllowlist
+   * rather than being handled here.
+   *
+   * The walk does not descend into a ReusedExchangeExec (a leaf whose wrapped 
exchange is a field,
+   * not a tree child), so a REUSED shuffle exchange would keep 
pipelined=false while its standalone
+   * twin flips to true. That divergence is not reachable: a reused shuffle 
requires
+   * referencing the same streaming source more than once (self-join / 
self-union / CTE read twice),
+   * which Real-Time Mode rejects when the query starts (MicroBatchExecution,
+   * IDENTICAL_SOURCES_IN_UNION_NOT_SUPPORTED) before this rule runs. The only 
ReusedExchangeExec
+   * that reaches an RTM plan wraps a BROADCAST exchange (multiple broadcast 
joins on the same
+   * static table, SC-209926), which this rule does not match.
+   *
+   * A pipelined shuffle read by more than one consumer (fan-out) is rejected 
by the DAGScheduler
+   * (checkPipelinedGroupsSupportedInRDDGraph). Note that check runs in 
handleJobSubmitted, so it
+   * rejects the batch's job rather than the query: such a query fails the 
same way on every batch
+   * instead of failing once when it is planned. Marking is not what makes the 
shape unsupported, so

Review Comment:
   fair will reword



-- 
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