Skip to content

SPARK-1202 - Add a "cancel" button in the UI for stages - #246

Closed
sundeepn wants to merge 4 commits into
apache:masterfrom
sundeepn:uikilljob
Closed

sundeepn wants to merge 4 commits into
apache:masterfrom
sundeepn:uikilljob

Conversation

@sundeepn

Copy link
Copy Markdown
Contributor

No description provided.

@AmplabJenkins

Copy link
Copy Markdown

Build triggered.

@AmplabJenkins

Copy link
Copy Markdown

Build started.

@AmplabJenkins

Copy link
Copy Markdown

Build finished.

@AmplabJenkins

Copy link
Copy Markdown

One or more automated tests failed
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/13489/

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Space after "running"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

As long as we're being nitpicky... should be "it's"

@andrewor14

Copy link
Copy Markdown
Contributor

Hi @sundeepn. Thanks for doing this. I left a few relatively minor comments. I haven't looked too closely into it but could you explain on a high level, what pages link to the KillPage and vice versa? If I understand correctly, each KillPage is associated with a stage ID, which seems a little strange when the "unit of killing" here is the job. Is there a reason why you created a new page for killing jobs, as opposed to embedding the functionality in existing pages (e.g. the StagePage)?

@sundeepn

Copy link
Copy Markdown
Contributor Author

Currently, only the IndexPage links to the killPage. The KillPage just has the uniform header on top so it links back to the other pages.

The unit of killing is a job by requirement. However, the initiation of kill action is through any stage contained in the job. The KillPage is actually a Job level message page, even though we arrive there by a stageId. Maybe I should change the title to "Job x containing stage y killed" ?

@AmplabJenkins

Copy link
Copy Markdown

Can one of the admins verify this patch?

@sundeepn

Copy link
Copy Markdown
Contributor Author

@andrewor14 / @kayousterhout. Thanks for the comments.
I have posted an additional commit with the changes. I have checked them in my branch but some reason they are not showing up here. Do you want me to issue another new pull request? Let me know.

sundeepn@32c2ea5

@AmplabJenkins

Copy link
Copy Markdown

Build triggered. Build is starting -or- tests failed to complete.

@AmplabJenkins

Copy link
Copy Markdown

Build started. Build is starting -or- tests failed to complete.

@AmplabJenkins

Copy link
Copy Markdown

Build finished. Build is starting -or- tests failed to complete.

@AmplabJenkins

Copy link
Copy Markdown

Build is starting -or- tests failed to complete.
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/13581/

@kayousterhout

Copy link
Copy Markdown
Contributor

It looks like github is just moving slowly today...the commit just got pulled in. I took another look at this and have a question: what happens for stages that are used for multiple jobs? Right now, stageIdToJobId in the UI code you added just maps a stage to a single job id. So, if stage0 is used by JobA and jobB, the ui code only stores one of these jobs, and then cancelJob() will only be called for one of the jobs. cancelJob() ultimately calls DAGScheduler.handleJobCancellation(), which only cancels the stages that are independent to the job. So, because stage0 is not independent to either of the jobs, it won't get cancelled. Did I misunderstand this?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

When we chatted about this I remember you saying that this code is to handle the case where the stage runs locally at the driver...but from glancing at the DAGScheduler code, it looks like onStageSubmitted() never gets called for the locally-run tasks, and they never show up at the UI.

When do you need this code / when will the jobIdToStageIds mapping not already be set up correctly by OnJobStart?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Well, this is only to ensure we can handle things if we get any scenarios where the onJobStart does not arrive before stageSubmitted. I am not familiar with the scheduling code sufficiently to rule that out. If you are sure, I can take this out.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Ah cool -- I looked at the ordering of the JobStart and StageSubmitted events more closely and I think you can safely remove this.

@sundeepn

Copy link
Copy Markdown
Contributor Author

what happens for stages that are used for multiple jobs?

I misunderstood our conversation on job to stage mapping the other day. As you see, the code will currently not handle multiple job mappings. Is there a simple example I can use to generate such a scenario?

@kayousterhout

Copy link
Copy Markdown
Contributor

I'm not sure how to generate an example of this...I think it can happen with Shark but maybe @pwendell or @markhamstra can comment here?

I wonder if maybe a better way to handle this would be to add a cancelStage() method to the DAGScheduler, which already has the mappings of stages to jobs and vice versa. That way the UI wouldn't have to duplicate this information.

@markhamstra

Copy link
Copy Markdown
Contributor

You can look at Matei's comments from way back: mesos/spark#414

@sundeepn

Copy link
Copy Markdown
Contributor Author

@kayousterhout I have thought about adding cancelStage to SparkContext/DAGScheduler. My earlier take was that the Job level Info is useful in the UI as well. Currently, its not shown/used in the UI, but lots of our users get confused looking at the stage level UI showing their query multiple times. It will be good to expose the Job ID in the stage table as well.

Having said that, I do not see any other way of doing this. We will have to move this to the DAGScheduler as a cancelStage. If two jobs can share stages, then any handling from the UI side can have a race condition in the cancel work flow and it will be a lot cleaner to handle upstream. I will submit a revision shortly.

@andrewor14, What do you think about adding a jobId column to the Stage table?

Thanks Mark for the pointer.

@markhamstra

Copy link
Copy Markdown
Contributor

+1 on job-level UI, including job-level progress indication.

@andrewor14

Copy link
Copy Markdown
Contributor

Yes, right now the landing page of the UI is the stage page, which is a little arbitrary. I think it makes sense to have an overview page, that displays the number of running executors, master URL, total duration, event log location and other more general info. (@tdas and I discussed this a bunch in designing for a brother UI for Spark Streaming.)

On this overview page, maybe we can also add a Job summary table, in which each row links to a job-specific stages page (what we already have). Then, in the context of canceling jobs, we can have the cancel button at the end of each Job row.

@pwendell

Copy link
Copy Markdown
Contributor

Jenkins, retest this please.

@AmplabJenkins

Copy link
Copy Markdown

Merged build triggered.

@AmplabJenkins

Copy link
Copy Markdown

Merged build started.

@AmplabJenkins

Copy link
Copy Markdown

Merged build finished. All automated tests passed.

@AmplabJenkins

Copy link
Copy Markdown

All automated tests passed.
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/14026/

@pwendell

Copy link
Copy Markdown
Contributor

Hey I'm gonna go ahead and merge this - I'll have a follow-on patch with some changes... felt that was easier than doing another round-trip on the review.

@pwendell

Copy link
Copy Markdown
Contributor

Merged into master and 1.0.

@asfgit asfgit closed this in 2c55783 Apr 11, 2014
asfgit pushed a commit that referenced this pull request Apr 11, 2014
Author: Sundeep Narravula <sundeepn@superduel.local>
Author: Sundeep Narravula <sundeepn@dhcpx-204-110.corp.yahoo.com>

Closes #246 from sundeepn/uikilljob and squashes the following commits:

5fdd0e2 [Sundeep Narravula] Fix test string
f6fdff1 [Sundeep Narravula] Format fix; reduced line size to less than 100 chars
d1daeb9 [Sundeep Narravula] Incorporating review comments.
8d97923 [Sundeep Narravula] Ability to kill jobs thru the UI. This behavior can be turned on be settings the following variable: spark.ui.killEnabled=true (default=false) Adding DAGScheduler event StageCancelled and corresponding handlers. Added cancellation reason to handlers.
(cherry picked from commit 2c55783)

Signed-off-by: Patrick Wendell <pwendell@gmail.com>
jhartlaub referenced this pull request in jhartlaub/spark May 27, 2014
Add missing license headers

I found this when doing further audits on the 0.8.1 release candidate.
(cherry picked from commit 6169fe1)

Signed-off-by: Patrick Wendell <pwendell@gmail.com>
pdeyhim pushed a commit to pdeyhim/spark-1 that referenced this pull request Jun 25, 2014
Previously, when jobs were cancelled, not all of the state in the
DAGScheduler was cleaned up, leading to a slow memory leak in the
DAGScheduler.  As we expose easier ways to cancel jobs, it's more
important to fix these issues.

This commit also fixes a second and less serious problem, which is that
previously, when a stage failed, not all of the appropriate stages
were cancelled.  See the "failure of stage used by two jobs" test
for an example of this.  This just meant that extra work was done, and is
not a correctness problem.

This commit adds 3 tests.  “run shuffle with map stage failure” is
a new test to more thoroughly test this functionality, and passes on
both the old and new versions of the code.  “trivial job
cancellation” fails on the old code because all state wasn’t cleaned
up correctly when jobs were cancelled (we didn’t remove the job from
resultStageToJob).  “failure of stage used by two jobs” fails on the
old code because taskScheduler.cancelTasks wasn’t called for one of
the stages (see test comments).

This should be checked in before apache#246, which makes it easier to
cancel stages / jobs.

Author: Kay Ousterhout <kayousterhout@gmail.com>

Closes apache#305 from kayousterhout/incremental_abort_fix and squashes the following commits:

f33d844 [Kay Ousterhout] Mark review comments
9217080 [Kay Ousterhout] Properly cleanup DAGScheduler on job cancellation.
pdeyhim pushed a commit to pdeyhim/spark-1 that referenced this pull request Jun 25, 2014
Author: Sundeep Narravula <sundeepn@superduel.local>
Author: Sundeep Narravula <sundeepn@dhcpx-204-110.corp.yahoo.com>

Closes apache#246 from sundeepn/uikilljob and squashes the following commits:

5fdd0e2 [Sundeep Narravula] Fix test string
f6fdff1 [Sundeep Narravula] Format fix; reduced line size to less than 100 chars
d1daeb9 [Sundeep Narravula] Incorporating review comments.
8d97923 [Sundeep Narravula] Ability to kill jobs thru the UI. This behavior can be turned on be settings the following variable: spark.ui.killEnabled=true (default=false) Adding DAGScheduler event StageCancelled and corresponding handlers. Added cancellation reason to handlers.
davies pushed a commit to davies/spark that referenced this pull request Apr 14, 2015
asfgit pushed a commit that referenced this pull request Apr 17, 2015
This PR pulls in recent changes in SparkR-pkg, including

cartesian, intersection, sampleByKey, subtract, subtractByKey, except, and some API for StructType and StructField.

Author: cafreeman <cfreeman@alteryx.com>
Author: Davies Liu <davies@databricks.com>
Author: Zongheng Yang <zongheng.y@gmail.com>
Author: Shivaram Venkataraman <shivaram.venkataraman@gmail.com>
Author: Shivaram Venkataraman <shivaram@cs.berkeley.edu>
Author: Sun Rui <rui.sun@intel.com>

Closes #5436 from davies/R3 and squashes the following commits:

c2b09be [Davies Liu] SQLTypes -> schema
a5a02f2 [Davies Liu] Merge branch 'master' of github.com:apache/spark into R3
168b7fe [Davies Liu] sort generics
b1fe460 [Davies Liu] fix conflict in README.md
e74c04e [Davies Liu] fix schema.R
4f5ac09 [Davies Liu] Merge branch 'master' of github.com:apache/spark into R5
41f8184 [Davies Liu] rm man
ae78312 [Davies Liu] Merge pull request #237 from sun-rui/SPARKR-154_3
1bdcb63 [Zongheng Yang] Updates to README.md.
5a553e7 [cafreeman] Use object attribute instead of argument
71372d9 [cafreeman] Update docs and examples
8526d2e [cafreeman] Remove `tojson` functions
6ef5f2d [cafreeman] Fix spacing
7741d66 [cafreeman] Rename the SQL DataType function
141efd8 [Shivaram Venkataraman] Merge pull request #245 from hqzizania/upstream
9387402 [Davies Liu] fix style
40199eb [Shivaram Venkataraman] Move except into sorted position
07d0dbc [Sun Rui] [SPARKR-244] Fix test failure after integration of subtract() and subtractByKey() for RDD.
7e8caa3 [Shivaram Venkataraman] Merge pull request #246 from hlin09/fixCombineByKey
ed66c81 [cafreeman] Update `subtract` to work with `generics.R`
f3ba785 [cafreeman] Fixed duplicate export
275deb4 [cafreeman] Update `NAMESPACE` and tests
1a3b63d [cafreeman] new version of `CreateDF`
836c4bf [cafreeman] Update `createDataFrame` and `toDF`
be5d5c1 [cafreeman] refactor schema functions
40338a4 [Zongheng Yang] Merge pull request #244 from sun-rui/SPARKR-154_5
20b97a6 [Zongheng Yang] Merge pull request #234 from hqzizania/assist
ba54e34 [Shivaram Venkataraman] Merge pull request #238 from sun-rui/SPARKR-154_4
c9497a3 [Shivaram Venkataraman] Merge pull request #208 from lythesia/master
b317aa7 [Zongheng Yang] Merge pull request #243 from hqzizania/master
136a07e [Zongheng Yang] Merge pull request #242 from hqzizania/stats
cd66603 [cafreeman] new line at EOF
8b76e81 [Shivaram Venkataraman] Merge pull request #233 from redbaron/fail-early-on-missing-dep
7dd81b7 [cafreeman] Documentation
0e2a94f [cafreeman] Define functions for schema and fields
mccheah pushed a commit to mccheah/spark that referenced this pull request Oct 12, 2017
[NOSQUASH] Resync from apache-spark-on-k8s upstream
bzhaoopenstack pushed a commit to bzhaoopenstack/spark that referenced this pull request Sep 11, 2019
cloud-fan pushed a commit that referenced this pull request Nov 10, 2020
### What changes were proposed in this pull request?
Push down filter through expand.  For case below:
```
create table t1(pid int, uid int, sid int, dt date, suid int) using parquet;
create table t2(pid int, vs int, uid int, csid int) using parquet;

SELECT
       years,
       appversion,
       SUM(uusers) AS users
FROM   (SELECT
               Date_trunc('year', dt)          AS years,
               CASE
                 WHEN h.pid = 3 THEN 'iOS'
                 WHEN h.pid = 4 THEN 'Android'
                 ELSE 'Other'
               END                             AS viewport,
               h.vs                            AS appversion,
               Count(DISTINCT u.uid)           AS uusers
               ,Count(DISTINCT u.suid)         AS srcusers
        FROM   t1 u
               join t2 h
                 ON h.uid = u.uid
        GROUP  BY 1,
                  2,
                  3) AS a
WHERE  viewport = 'iOS'
GROUP  BY 1,
          2
```

Plan. before this pr:
```
== Physical Plan ==
*(5) HashAggregate(keys=[years#30, appversion#32], functions=[sum(uusers#33L)])
+- Exchange hashpartitioning(years#30, appversion#32, 200), true, [id=#251]
   +- *(4) HashAggregate(keys=[years#30, appversion#32], functions=[partial_sum(uusers#33L)])
      +- *(4) HashAggregate(keys=[date_trunc('year', CAST(u.`dt` AS TIMESTAMP))#45, CASE WHEN (h.`pid` = 3) THEN 'iOS' WHEN (h.`pid` = 4) THEN 'Android' ELSE 'Other' END#46, vs#12], functions=[count(if ((gid#44 = 1)) u.`uid`#47 else null)])
         +- Exchange hashpartitioning(date_trunc('year', CAST(u.`dt` AS TIMESTAMP))#45, CASE WHEN (h.`pid` = 3) THEN 'iOS' WHEN (h.`pid` = 4) THEN 'Android' ELSE 'Other' END#46, vs#12, 200), true, [id=#246]
            +- *(3) HashAggregate(keys=[date_trunc('year', CAST(u.`dt` AS TIMESTAMP))#45, CASE WHEN (h.`pid` = 3) THEN 'iOS' WHEN (h.`pid` = 4) THEN 'Android' ELSE 'Other' END#46, vs#12], functions=[partial_count(if ((gid#44 = 1)) u.`uid`#47 else null)])
               +- *(3) HashAggregate(keys=[date_trunc('year', CAST(u.`dt` AS TIMESTAMP))#45, CASE WHEN (h.`pid` = 3) THEN 'iOS' WHEN (h.`pid` = 4) THEN 'Android' ELSE 'Other' END#46, vs#12, u.`uid`#47, u.`suid`#48, gid#44], functions=[])
                  +- Exchange hashpartitioning(date_trunc('year', CAST(u.`dt` AS TIMESTAMP))#45, CASE WHEN (h.`pid` = 3) THEN 'iOS' WHEN (h.`pid` = 4) THEN 'Android' ELSE 'Other' END#46, vs#12, u.`uid`#47, u.`suid`#48, gid#44, 200), true, [id=#241]
                     +- *(2) HashAggregate(keys=[date_trunc('year', CAST(u.`dt` AS TIMESTAMP))#45, CASE WHEN (h.`pid` = 3) THEN 'iOS' WHEN (h.`pid` = 4) THEN 'Android' ELSE 'Other' END#46, vs#12, u.`uid`#47, u.`suid`#48, gid#44], functions=[])
                        +- *(2) Filter (CASE WHEN (h.`pid` = 3) THEN 'iOS' WHEN (h.`pid` = 4) THEN 'Android' ELSE 'Other' END#46 = iOS)
                           +- *(2) Expand [ArrayBuffer(date_trunc(year, cast(dt#9 as timestamp), Some(Etc/GMT+7)), CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END, vs#12, uid#7, null, 1), ArrayBuffer(date_trunc(year, cast(dt#9 as timestamp), Some(Etc/GMT+7)), CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END, vs#12, null, suid#10, 2)], [date_trunc('year', CAST(u.`dt` AS TIMESTAMP))#45, CASE WHEN (h.`pid` = 3) THEN 'iOS' WHEN (h.`pid` = 4) THEN 'Android' ELSE 'Other' END#46, vs#12, u.`uid`#47, u.`suid`#48, gid#44]
                              +- *(2) Project [uid#7, dt#9, suid#10, pid#11, vs#12]
                                 +- *(2) BroadcastHashJoin [uid#7], [uid#13], Inner, BuildRight
                                    :- *(2) Project [uid#7, dt#9, suid#10]
                                    :  +- *(2) Filter isnotnull(uid#7)
                                    :     +- *(2) ColumnarToRow
                                    :        +- FileScan parquet default.t1[uid#7,dt#9,suid#10] Batched: true, DataFilters: [isnotnull(uid#7)], Format: Parquet, Location: InMemoryFileIndex[file:/root/spark-3.0.0-bin-hadoop3.2/spark-warehouse/t1], PartitionFilters: [], PushedFilters: [IsNotNull(uid)], ReadSchema: struct<uid:int,dt:date,suid:int>
                                    +- BroadcastExchange HashedRelationBroadcastMode(List(cast(input[2, int, true] as bigint))), [id=#233]
                                       +- *(1) Project [pid#11, vs#12, uid#13]
                                          +- *(1) Filter isnotnull(uid#13)
                                             +- *(1) ColumnarToRow
                                                +- FileScan parquet default.t2[pid#11,vs#12,uid#13] Batched: true, DataFilters: [isnotnull(uid#13)], Format: Parquet, Location: InMemoryFileIndex[file:/root/spark-3.0.0-bin-hadoop3.2/spark-warehouse/t2], PartitionFilters: [], PushedFilters: [IsNotNull(uid)], ReadSchema: struct<pid:int,vs:int,uid:int>
```

Plan. after. this pr. :
```
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- HashAggregate(keys=[years#0, appversion#2], functions=[sum(uusers#3L)], output=[years#0, appversion#2, users#5L])
   +- Exchange hashpartitioning(years#0, appversion#2, 5), true, [id=#71]
      +- HashAggregate(keys=[years#0, appversion#2], functions=[partial_sum(uusers#3L)], output=[years#0, appversion#2, sum#22L])
         +- HashAggregate(keys=[date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles))#23, CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END#24, vs#12], functions=[count(distinct uid#7)], output=[years#0, appversion#2, uusers#3L])
            +- Exchange hashpartitioning(date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles))#23, CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END#24, vs#12, 5), true, [id=#67]
               +- HashAggregate(keys=[date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles))#23, CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END#24, vs#12], functions=[partial_count(distinct uid#7)], output=[date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles))#23, CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END#24, vs#12, count#27L])
                  +- HashAggregate(keys=[date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles))#23, CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END#24, vs#12, uid#7], functions=[], output=[date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles))#23, CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END#24, vs#12, uid#7])
                     +- Exchange hashpartitioning(date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles))#23, CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END#24, vs#12, uid#7, 5), true, [id=#63]
                        +- HashAggregate(keys=[date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles)) AS date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles))#23, CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END AS CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END#24, vs#12, uid#7], functions=[], output=[date_trunc(year, cast(dt#9 as timestamp), Some(America/Los_Angeles))#23, CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END#24, vs#12, uid#7])
                           +- Project [uid#7, dt#9, pid#11, vs#12]
                              +- BroadcastHashJoin [uid#7], [uid#13], Inner, BuildRight, false
                                 :- Filter isnotnull(uid#7)
                                 :  +- FileScan parquet default.t1[uid#7,dt#9] Batched: true, DataFilters: [isnotnull(uid#7)], Format: Parquet, Location: InMemoryFileIndex[file:/private/var/folders/4l/7_c5c97s1_gb0d9_d6shygx00000gn/T/warehouse-c069d87..., PartitionFilters: [], PushedFilters: [IsNotNull(uid)], ReadSchema: struct<uid:int,dt:date>
                                 +- BroadcastExchange HashedRelationBroadcastMode(List(cast(input[2, int, false] as bigint)),false), [id=#58]
                                    +- Filter ((CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END = iOS) AND isnotnull(uid#13))
                                       +- FileScan parquet default.t2[pid#11,vs#12,uid#13] Batched: true, DataFilters: [(CASE WHEN (pid#11 = 3) THEN iOS WHEN (pid#11 = 4) THEN Android ELSE Other END = iOS), isnotnull..., Format: Parquet, Location: InMemoryFileIndex[file:/private/var/folders/4l/7_c5c97s1_gb0d9_d6shygx00000gn/T/warehouse-c069d87..., PartitionFilters: [], PushedFilters: [IsNotNull(uid)], ReadSchema: struct<pid:int,vs:int,uid:int>

```

### Why are the changes needed?
Improve  performance, filter more data.

### Does this PR introduce _any_ user-facing change?
No

### How was this patch tested?
Added UT

Closes #30278 from AngersZhuuuu/SPARK-33302.

Authored-by: angerszhu <angers.zhu@gmail.com>
Signed-off-by: Wenchen Fan <wenchen@databricks.com>
MaxGekk added a commit to MaxGekk/spark that referenced this pull request Sep 18, 2026
### What changes were proposed in this pull request?

`PLAN_TASK_104.md`, the plan for `Long` arithmetic - task 30's int64 half - and the milestone row updated with what planning it found. No code.

**Reading the emitter before planning changed what this task is.** The row reads as though the long-lane arithmetic has to be built. It largely does not. `emitIntArith` and `emitIntNeg` were converted to `Lane` descriptors by task 85 step 3 and carry no int32 assumption: the lanewise op, the two XORs, the AND, the sign compare and the broadcast all go through `analysis.lane.*`. `VarkaLaneTypeSuite` already asserts both nodes are emittable at `LONG`, and `VarkaReferenceEvaluator.evalLong` already has arms for them. So checked add, subtract and negate at the long lane are already emitted correctly, and what stops `bi + 1` fusing is the **compiler**, which task 29 deliberately narrowed to comparisons, `greatest`/`least` and `CASE WHEN`.

**What is actually missing is the lattice, and it is the thing that removes the checks.** `VarkaRangeAnalysis.range` opens by answering `UNKNOWN` for every node that is not on the int lane, with its reasons stated: the epoch-day contract is an int32 statement, and `literals` hands back an int, which cannot carry a 64-bit constant. That refusal is correct and it is also this task's blocker, because `cannotOverflow` is the only thing that takes a check off. At the long lane it can never be true today, so every checked long op would keep its four-op test even where a literal operand makes overflow impossible, and every checked long multiply would decline unconditionally.

The plan owes two concrete things there, both named in the refusal itself: a long literal table the analysis can read (`magnitude` takes an `IntUnaryOperator`), and a saturation rule that is lane-dependent rather than int32 - `cannotOverflow` compares a bound against `Int.MaxValue`, and widening that comparison without making it lane-aware would take the check off a long op that can overflow, which is a wrong-answer bug rather than a performance one. There is no epoch-day contract at this lane and the plan invents none: a `bigint` column is unbounded, a `TIME` column is bounded by its type, a day-time interval is not. So the lattice's long half starts as "literals and the types that bound themselves", which takes the check off `bi + 1` and not off `bi + bi2`, and the tests pin that asymmetry rather than hide it.

**The checked multiply is the one real design question, and it has two halves.** At int32 the recorded fix (`PLAN_TASK_63.md` 3.4, quoted in the compiler today) is exactly task 28 plus this lane: widen both operands, multiply in int64 where `2^31 * 2^31 = 2^62` cannot overflow, range-check the result. That closes task 30's last int32 gap. At **int64 there is no wider lane** and `VectorOperators` has no multiply-high, so the plan lays out three candidates to build and measure rather than choosing by argument: bound-only (ships first, zero ops when proved), a double-lane magnitude pre-check, and a 64x64 high-product emulation from 32-bit halves.

The middle one is worth its paragraph because its soundness is not obvious: converting a long above 2^53 to double loses precision, but the loss is *relative*, so the product estimate is within about 2^-52 of the true magnitude; testing against 2^62 when the real limit is 2^63 leaves a factor of two of margin. It can therefore never admit an overflow, only decline a product that would have been fine - and declining is always safe. It also carries an AVX-level caveat the other two do not: `L2DProbe` records that under `-XX:UseAVX=2` the long-to-double converts fail to inline and the loop goes scalar, so it may be a good default at AVX-512 and a bad one at 128-bit.

**Dependencies, checked rather than assumed.** Section 2.39 ends "and `Cast` between int and long as 2.2's conversion", and the wave table lists 104 after 29, 84 and 28's cast. 29 and 84 are merged. Task 28 is planned (apache#246) and not built, and it owns that cast - this task consumes it. Everything that is this task's own subject is single-lane and unblocked today.

### Why are the changes needed?

Task 30's row closes when this one does, and the milestone's own sequencing puts 104 in wave 4 with task 28. Planning it first is what surfaced the two findings above: that most of the emitter work is already done, and that the lattice - not the emitter - is what the task actually has to build.

### Does this PR introduce _any_ user-facing change?

No. A plan and a milestone row.

### How was this patch tested?

`dev/varka_precommit.sh` and `dev/varka_quote_check.py` green; the non-ASCII scan clean. No code, so no suites apply.

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

Generated-by: Claude Code (Claude Opus 5)
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.

6 participants