Commit 3d0058f
committed
[SPARK-XXXXX][CORE] Add ConcurrentStageDAGScheduler for low-latency streaming
Ports the ConcurrentStageDAGScheduler from the Databricks runtime so that
streaming queries can opt in to a "real-time" execution mode that runs all
stages of a job concurrently rather than sequentially.
When enabled via spark.scheduler.dagSchedulerType=ConcurrentStageDAGScheduler
and the per-job streaming.concurrent.stages.enabled property, the scheduler:
- Marks all ancestor stages of the final stage as concurrent on job submission
and validates that the cluster has enough free slots
(CONCURRENT_SCHEDULER_INSUFFICIENT_SLOT), gated by
spark.scheduler.realtimeModeSlotsCheck.disabled.
- Submits child stages while parents are still running, delays task completion
events for a child whose parent is still running, and replays the delayed
events when the parent finishes.
- Rejects speculative execution.
DAGScheduler changes (no-op for the default scheduler):
- New protected onFinalStageCreated hook, invoked from handleJobSubmitted /
handleMapStageSubmitted right after final stage creation.
- New protected submitConcurrentStage and postSchedulerEvent helpers.
- New package-visible isRunningStage and getStage accessors.
- submitStage and markStageAsFinished relaxed from private to protected so
subclasses can override them.
DAGSchedulerSuite refactor:
- Renames the concrete suite to abstract DAGSchedulerSuiteBase and adds an
empty class DAGSchedulerSuite extends DAGSchedulerSuiteBase to preserve
the existing entry point.
- Extracts a TestDAGScheduler trait carrying the scheduleShuffleMergeFinalize
and handleTaskCompletion overrides; MyDAGScheduler mixes the trait in.
- Adds a protected createInitialScheduler hook used by init().
- Loosens submit, completeShuffleMapStageSuccessfully,
completeNextResultStageWithSuccess, and assertDataStructuresEmpty to
protected so subclass suites can use them.
Integration:
- SparkContext picks the scheduler implementation based on
spark.scheduler.dagSchedulerType.
- TaskSchedulerImpl uses maxFailures=1 for concurrent-stage TaskSets so a
failure restarts the streaming query instead of being silently retried.
- TaskSetManager counts ExecutorLostFailure toward task failures and skips
the "executor lost is not the task's fault" exemption in concurrent mode.
Adds the supporting LogKeys (PARENT_STAGE, STREAMING_QUERY_ID) and the
CONCURRENT_SCHEDULER_INSUFFICIENT_SLOT error class.
Deviations from the runtime source kept to the minimum necessary to compile
in OSS:
- Extends DAGScheduler directly (runtime extends CrossJobDepDAGScheduler,
which gates micro-batch pipelining; not part of OSS).
- Hook is named onFinalStageCreated rather than the runtime's
populateCrossJobDepInfo, since CrossJobDepDAGScheduler is not part of OSS.
- Micro-batch pipelining co-existence check (and its test) dropped, since
MBP is not part of OSS.
- getStreamingBatchIdFromProperties and StreamingBatchId live in the
companion object instead of CrossJobDepDAGScheduler.
- Slot check uses sc.schedulerBackend.defaultParallelism() in place of
the runtime's TaskSchedulerStats helper.
- DatabricksEdgeConfigs.serverlessEnabled gating removed; the
spark.scheduler.realtimeModeSlotsCheck.disabled config is the sole knob.
- isConcurrentStagesEnabled tolerates null Properties (OSS TaskSet allows
null in tests).
Co-authored-by: Isaac1 parent 706b6a3 commit 3d0058f
12 files changed
Lines changed: 774 additions & 32 deletions
File tree
- common
- utils-java/src/main/java/org/apache/spark/internal
- utils/src/main/resources/error
- core/src
- main/scala/org/apache/spark
- internal/config
- scheduler
- test/scala/org/apache/spark/scheduler
Lines changed: 2 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
577 | 577 | | |
578 | 578 | | |
579 | 579 | | |
| 580 | + | |
580 | 581 | | |
581 | 582 | | |
582 | 583 | | |
| |||
792 | 793 | | |
793 | 794 | | |
794 | 795 | | |
| 796 | + | |
795 | 797 | | |
796 | 798 | | |
797 | 799 | | |
| |||
Lines changed: 6 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
890 | 890 | | |
891 | 891 | | |
892 | 892 | | |
| 893 | + | |
| 894 | + | |
| 895 | + | |
| 896 | + | |
| 897 | + | |
| 898 | + | |
893 | 899 | | |
894 | 900 | | |
895 | 901 | | |
| |||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
600 | 600 | | |
601 | 601 | | |
602 | 602 | | |
603 | | - | |
| 603 | + | |
| 604 | + | |
| 605 | + | |
| 606 | + | |
| 607 | + | |
604 | 608 | | |
605 | 609 | | |
606 | 610 | | |
| |||
Lines changed: 18 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
2396 | 2396 | | |
2397 | 2397 | | |
2398 | 2398 | | |
| 2399 | + | |
| 2400 | + | |
| 2401 | + | |
| 2402 | + | |
| 2403 | + | |
| 2404 | + | |
| 2405 | + | |
| 2406 | + | |
| 2407 | + | |
| 2408 | + | |
| 2409 | + | |
| 2410 | + | |
| 2411 | + | |
| 2412 | + | |
| 2413 | + | |
| 2414 | + | |
| 2415 | + | |
| 2416 | + | |
2399 | 2417 | | |
2400 | 2418 | | |
2401 | 2419 | | |
| |||
Lines changed: 282 additions & 0 deletions
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
| |||
| 1 | + | |
| 2 | + | |
| 3 | + | |
| 4 | + | |
| 5 | + | |
| 6 | + | |
| 7 | + | |
| 8 | + | |
| 9 | + | |
| 10 | + | |
| 11 | + | |
| 12 | + | |
| 13 | + | |
| 14 | + | |
| 15 | + | |
| 16 | + | |
| 17 | + | |
| 18 | + | |
| 19 | + | |
| 20 | + | |
| 21 | + | |
| 22 | + | |
| 23 | + | |
| 24 | + | |
| 25 | + | |
| 26 | + | |
| 27 | + | |
| 28 | + | |
| 29 | + | |
| 30 | + | |
| 31 | + | |
| 32 | + | |
| 33 | + | |
| 34 | + | |
| 35 | + | |
| 36 | + | |
| 37 | + | |
| 38 | + | |
| 39 | + | |
| 40 | + | |
| 41 | + | |
| 42 | + | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
| 48 | + | |
| 49 | + | |
| 50 | + | |
| 51 | + | |
| 52 | + | |
| 53 | + | |
| 54 | + | |
| 55 | + | |
| 56 | + | |
| 57 | + | |
| 58 | + | |
| 59 | + | |
| 60 | + | |
| 61 | + | |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
| 76 | + | |
| 77 | + | |
| 78 | + | |
| 79 | + | |
| 80 | + | |
| 81 | + | |
| 82 | + | |
| 83 | + | |
| 84 | + | |
| 85 | + | |
| 86 | + | |
| 87 | + | |
| 88 | + | |
| 89 | + | |
| 90 | + | |
| 91 | + | |
| 92 | + | |
| 93 | + | |
| 94 | + | |
| 95 | + | |
| 96 | + | |
| 97 | + | |
| 98 | + | |
| 99 | + | |
| 100 | + | |
| 101 | + | |
| 102 | + | |
| 103 | + | |
| 104 | + | |
| 105 | + | |
| 106 | + | |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
| 112 | + | |
| 113 | + | |
| 114 | + | |
| 115 | + | |
| 116 | + | |
| 117 | + | |
| 118 | + | |
| 119 | + | |
| 120 | + | |
| 121 | + | |
| 122 | + | |
| 123 | + | |
| 124 | + | |
| 125 | + | |
| 126 | + | |
| 127 | + | |
| 128 | + | |
| 129 | + | |
| 130 | + | |
| 131 | + | |
| 132 | + | |
| 133 | + | |
| 134 | + | |
| 135 | + | |
| 136 | + | |
| 137 | + | |
| 138 | + | |
| 139 | + | |
| 140 | + | |
| 141 | + | |
| 142 | + | |
| 143 | + | |
| 144 | + | |
| 145 | + | |
| 146 | + | |
| 147 | + | |
| 148 | + | |
| 149 | + | |
| 150 | + | |
| 151 | + | |
| 152 | + | |
| 153 | + | |
| 154 | + | |
| 155 | + | |
| 156 | + | |
| 157 | + | |
| 158 | + | |
| 159 | + | |
| 160 | + | |
| 161 | + | |
| 162 | + | |
| 163 | + | |
| 164 | + | |
| 165 | + | |
| 166 | + | |
| 167 | + | |
| 168 | + | |
| 169 | + | |
| 170 | + | |
| 171 | + | |
| 172 | + | |
| 173 | + | |
| 174 | + | |
| 175 | + | |
| 176 | + | |
| 177 | + | |
| 178 | + | |
| 179 | + | |
| 180 | + | |
| 181 | + | |
| 182 | + | |
| 183 | + | |
| 184 | + | |
| 185 | + | |
| 186 | + | |
| 187 | + | |
| 188 | + | |
| 189 | + | |
| 190 | + | |
| 191 | + | |
| 192 | + | |
| 193 | + | |
| 194 | + | |
| 195 | + | |
| 196 | + | |
| 197 | + | |
| 198 | + | |
| 199 | + | |
| 200 | + | |
| 201 | + | |
| 202 | + | |
| 203 | + | |
| 204 | + | |
| 205 | + | |
| 206 | + | |
| 207 | + | |
| 208 | + | |
| 209 | + | |
| 210 | + | |
| 211 | + | |
| 212 | + | |
| 213 | + | |
| 214 | + | |
| 215 | + | |
| 216 | + | |
| 217 | + | |
| 218 | + | |
| 219 | + | |
| 220 | + | |
| 221 | + | |
| 222 | + | |
| 223 | + | |
| 224 | + | |
| 225 | + | |
| 226 | + | |
| 227 | + | |
| 228 | + | |
| 229 | + | |
| 230 | + | |
| 231 | + | |
| 232 | + | |
| 233 | + | |
| 234 | + | |
| 235 | + | |
| 236 | + | |
| 237 | + | |
| 238 | + | |
| 239 | + | |
| 240 | + | |
| 241 | + | |
| 242 | + | |
| 243 | + | |
| 244 | + | |
| 245 | + | |
| 246 | + | |
| 247 | + | |
| 248 | + | |
| 249 | + | |
| 250 | + | |
| 251 | + | |
| 252 | + | |
| 253 | + | |
| 254 | + | |
| 255 | + | |
| 256 | + | |
| 257 | + | |
| 258 | + | |
| 259 | + | |
| 260 | + | |
| 261 | + | |
| 262 | + | |
| 263 | + | |
| 264 | + | |
| 265 | + | |
| 266 | + | |
| 267 | + | |
| 268 | + | |
| 269 | + | |
| 270 | + | |
| 271 | + | |
| 272 | + | |
| 273 | + | |
| 274 | + | |
| 275 | + | |
| 276 | + | |
| 277 | + | |
| 278 | + | |
| 279 | + | |
| 280 | + | |
| 281 | + | |
| 282 | + | |
0 commit comments