Skip to content

[SPARK-59924][K8s][DOCS] Fix stale "Valid values" list for spark.kubernetes.executor.rollPolicy - #59187

Closed
sarutak wants to merge 3 commits into
apache:masterfrom
sarutak:rollpolicy-valid-values-docfix
Closed

sarutak wants to merge 3 commits into
apache:masterfrom
sarutak:rollpolicy-valid-values-docfix

Conversation

@sarutak

@sarutak sarutak commented Oct 1, 2026 •

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

This PR fixes the "Valid values" list for spark.kubernetes.executor.rollPolicy, which was out of sync with the ExecutorRollPolicy enum in both Config.scala and running-on-kubernetes.md.

  • Config.scala: added the missing TOTAL_SHUFFLE_WRITE and DISK_USED.
  • running-on-kubernetes.md: added the missing AVERAGE_DURATION, PEAK_JVM_ONHEAP_MEMORY, PEAK_JVM_OFFHEAP_MEMORY, TOTAL_SHUFFLE_WRITE, DISK_USED, ACTIVE_TASKS and OUTLIER_NO_FALLBACK, including their per-policy descriptions.

Both lists now follow the enum definition order.

Why are the changes needed?

The documented list of valid values listed only a subset of the supported policies, which is misleading.

Does this PR introduce any user-facing change?

No. Docs only. No behavior change.

How was this patch tested?

Verified by inspection that the lists in both files match the ExecutorRollPolicy enum, and that the OUTLIER dimensions, their checking order, and the DISK_USED wording match ExecutorRollPlugin and AppStatusListener.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Kiro CLI / Claude

The "Valid values" summary for spark.kubernetes.executor.rollPolicy in
both Config.scala and running-on-kubernetes.md was out of sync with the
ExecutorRollPolicy enum.

- Config.scala omitted TOTAL_SHUFFLE_WRITE and DISK_USED from the
  "Valid values" list, even though both are valid policies and are
  described later in the same doc string.
- running-on-kubernetes.md omitted AVERAGE_DURATION, PEAK_JVM_ONHEAP_MEMORY,
  PEAK_JVM_OFFHEAP_MEMORY, TOTAL_SHUFFLE_WRITE, DISK_USED, ACTIVE_TASKS and
  OUTLIER_NO_FALLBACK, and lacked per-policy descriptions for several of them.

Both lists now match the enum definition order, and the per-policy
descriptions are consistent between the two files.

Also fix two missing spaces at string-concatenation boundaries in the
Config.scala doc string that rendered as "summary.ID" and "thanat".

Docs only, no behavior change (checkValues already validates against the
enum).

@HyukjinKwon HyukjinKwon left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Review summary

The valid-values lists in both Config.scala and running-on-kubernetes.md now match the 13 ExecutorRollPolicy values in enum order. Each new per-policy sentence matches the corresponding sort in ExecutorRollPlugin.choose. The only remaining gap is that the OUTLIER description, which the newly documented OUTLIER_NO_FALLBACK entry relies on, still names four of the eight outlier dimensions the plugin checks (inline comment). This is a small, non-blocking doc fix.

Findings

1 total: 0 P0, 0 P1, 0 P2, 1 P3.

Nit (P3)

  • OUTLIER criteria still list only 4 of the 8 checked dimensions — docs/running-on-kubernetes.md:2127 — see inline.

at least two standard deviation from the mean in average task time,
total task time, total task GC time, and the number of failed tasks if exists.
If there is no outlier, it works like TOTAL_DURATION policy.
OUTLIER_NO_FALLBACK policy picks an outlier using the OUTLIER policy above.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit (P3): This entry defines OUTLIER_NO_FALLBACK by pointing to the OUTLIER description above. That description (here and in the Config.scala doc string) still lists only average task time, total task time, total task GC time, and failed tasks. ExecutorRollPlugin.outliersFromMultipleDimensions also checks peak JVM on-heap memory, peak JVM off-heap memory, total shuffle write, and disk used. As written, users would expect an executor that is an outlier only in, say, disk used never to be rolled under OUTLIER_NO_FALLBACK, but it will be. Since this PR is syncing these docs with the enum, could you extend the OUTLIER dimension list in both files to all eight metrics?

Verification:

  • Inspection: The documented OUTLIER dimensions in both files match the outliers(...) calls in ExecutorRollPlugin.outliersFromMultipleDimensions.
  • Compile: The kubernetes core module still compiles, and the Scala source lines stay within the 100-character limit.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, makes sense. Extended the OUTLIER dimension list in both files to all eight metrics checked by outliersFromMultipleDimensions.

Address review feedback: the OUTLIER description listed only four of the
eight dimensions that ExecutorRollPlugin.outliersFromMultipleDimensions
checks. Add peak JVM on-heap memory, peak JVM off-heap memory, total
shuffle write, and disk used so OUTLIER / OUTLIER_NO_FALLBACK docs match
the implementation.

@dongjoon-hyun dongjoon-hyun left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you, @sarutak .

Review summary

The Valid values lists in both Config.scala and running-on-kubernetes.md now match the 13 ExecutorRollPolicy values, and each new per-policy sentence matches the corresponding sort in ExecutorRollPlugin.choose. The two missing-space fixes in the Config.scala doc string are correct. The OUTLIER dimension list was already extended after the earlier review.

I have two small, non-blocking doc suggestions (inline).

Findings

2 total: 0 P0, 0 P1, 0 P2, 2 P3.

Nits (P3)

  • DISK_USED description is ambiguous (docs/running-on-kubernetes.md:2118)
  • OUTLIER dimensions are checked in a priority order that is not documented (docs/running-on-kubernetes.md:2127)

Comment thread docs/running-on-kubernetes.md Outdated
PEAK_JVM_ONHEAP_MEMORY policy chooses an executor with the biggest peak JVM on-heap memory.
PEAK_JVM_OFFHEAP_MEMORY policy chooses an executor with the biggest peak JVM off-heap memory.
TOTAL_SHUFFLE_WRITE policy chooses an executor with the biggest total shuffle write.
DISK_USED policy chooses an executor with the biggest used disk size.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit (P3): "the biggest used disk size" can be read as the executor's local disk usage (scratch space or shuffle files). The metric behind DISK_USED is ExecutorSummary.diskUsed, which AppStatusListener accumulates only from SparkListenerBlockUpdated events for RDD, stream and broadcast blocks (updateExecutorMemoryDiskInfo). Shuffle files are not counted. An executor that holds large shuffle data but no disk-persisted blocks therefore never ranks first under this policy, and the same applies to the new "disk used" dimension of OUTLIER.

Suggestion: "DISK_USED policy chooses an executor with the biggest disk size used by its stored blocks (e.g., disk-persisted RDD blocks)." The same sentence is in the Config.scala doc string (line 264), which is outside this diff's hunks, so please update both.

total task time, total task GC time, and the number of failed tasks if exists.
total task time, total task GC time, the number of failed tasks,
peak JVM on-heap memory, peak JVM off-heap memory, total shuffle write,
and disk used if exists.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit (P3): This lists eight dimensions, but outliersFromMultipleDimensions concatenates the per-dimension outlier lists in a fixed order (average task time > total task time > total task GC time > failed tasks > peak JVM on-heap > peak JVM off-heap > total shuffle write > disk used) and choose takes the first element. When different executors are outliers in different dimensions, the documented list alone does not tell the reader which one is rolled. Since the list is written in that same order, one sentence such as "The dimensions are checked in this order and the first outlier found is chosen." would make it explicit. The same applies to the Config.scala doc string.

Address review feedback:
- DISK_USED is backed by ExecutorSummary.diskUsed, which AppStatusListener
  accumulates only from RDD/stream/broadcast block updates (not shuffle
  files). Clarify that it is the disk size used by stored blocks.
- outliersFromMultipleDimensions checks the eight dimensions in a fixed
  priority order and choose() takes the first outlier found. Document that
  order explicitly so readers know which executor is rolled.

Updated both Config.scala and running-on-kubernetes.md.
@sarutak

sarutak commented Oct 2, 2026

Copy link
Copy Markdown
Member Author

Thanks @dongjoon-hyun. I've addressed your comments:

  • Clarified that DISK_USED reflects the disk size used by stored blocks (RDD/stream/broadcast), not shuffle files, since that is what ExecutorSummary.diskUsed accumulates.
  • Documented that the OUTLIER dimensions are checked in the listed priority order and the first outlier found is chosen.

Updated both running-on-kubernetes.md and the Config.scala doc string.

@dongjoon-hyun

Copy link
Copy Markdown
Member

Thank you, @sarutak .

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+1, LGTM.

dongjoon-hyun pushed a commit that referenced this pull request Oct 2, 2026
…rnetes.executor.rollPolicy`

### What changes were proposed in this pull request?
This PR fixes the "Valid values" list for `spark.kubernetes.executor.rollPolicy`, which was out of sync with the `ExecutorRollPolicy` enum in both `Config.scala` and `running-on-kubernetes.md`.

- `Config.scala`: added the missing `TOTAL_SHUFFLE_WRITE` and `DISK_USED`.
- `running-on-kubernetes.md`: added the missing `AVERAGE_DURATION`, `PEAK_JVM_ONHEAP_MEMORY`, `PEAK_JVM_OFFHEAP_MEMORY`, `TOTAL_SHUFFLE_WRITE`, `DISK_USED`, `ACTIVE_TASKS` and `OUTLIER_NO_FALLBACK`, including their per-policy descriptions.

Both lists now follow the enum definition order.

### Why are the changes needed?
The documented list of valid values listed only a subset of the supported policies, which is misleading.

### Does this PR introduce _any_ user-facing change?
No. Docs only. No behavior change.

### How was this patch tested?
Verified by inspection that the lists in both files match the `ExecutorRollPolicy` enum, and that the `OUTLIER` dimensions, their checking order, and the `DISK_USED` wording match `ExecutorRollPlugin` and `AppStatusListener`.

### Was this patch authored or co-authored using generative AI tooling?
Generated-by: Kiro CLI / Claude

Closes #59187 from sarutak/rollpolicy-valid-values-docfix.

Authored-by: Kousuke Saruta <sarutak@apache.org>
Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>
(cherry picked from commit 25a9aa7)
Signed-off-by: Dongjoon Hyun <dongjoon@apache.org>
@dongjoon-hyun

Copy link
Copy Markdown
Member

Merge Summary:

Posted by merge_spark_pr.py

@dongjoon-hyun

Copy link
Copy Markdown
Member

To @sarutak , we need new backporting PRs from branch-4.3 because the following is Apache Spark 4.4 feature. If you want to backport, please make new PRs.

@sarutak

sarutak commented Oct 3, 2026

Copy link
Copy Markdown
Member Author

Thank you, @dongjoon-hyun and @HyukjinKwon !

To @sarutak , we need new backporting PRs from branch-4.3 because the following is Apache Spark 4.4 feature. If you want to backport, please make new PRs.

Got it. Will open a PR for branch-4.3.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants