Repository navigation
Conversation
|
Can one of the admins verify this patch? |
|
Jenkins, this is ok to test |
|
Merged build triggered. |
|
Merged build started. |
|
Merged build finished. |
|
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/14261/ |
|
Mind fixing the tab characters so we can test this? |
|
sure, I am looking at right now. |
|
Merged build triggered. |
|
Merged build started. |
|
Merged build finished. |
|
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/14295/ |
|
BTW, you can check style locally |
|
Merged build triggered. |
|
Merged build started. |
|
Merged build finished. |
|
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/14300/ |
|
The problem was I used vim for coding and it screwed up the tabbing for some reason. I did the sbt scalastyle and it succeed now locally |
|
I am a little bit confused why it does not complain in my machine locally but it produces errors here ... |
|
Merged build triggered. |
|
Merged build started. |
|
Merged build finished. |
|
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/14304/ |
|
Build triggered. |
|
Build started. |
|
Build finished. |
|
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/14496/ |
|
Build triggered. |
|
Build started. |
|
Build finished. |
|
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/14497/ |
## Upstream SPARK-XXXXX ticket and PR link (if not applicable, explain) Not filed in upstream, touches code for conda. ## What changes were proposed in this pull request? rLibDir contains a sequence of possible paths for the SparkR package on the executor and is passed on to the R daemon with the SPARKR_RLIBDIR environment variable. This PR filters rLibDir for paths that exist before setting SPARKR_RLIBDIR, retaining existing functionality to preferentially choose a YARN or local SparkR install over conda if both are present. See daemon.R: https://github.com/palantir/spark/blob/master/R/pkg/inst/worker/daemon.R#L23 Fixes apache#456 ## How was this patch tested? Manually testing cherry picked on older version Please review http://spark.apache.org/contributing.html before opening a pull request.
It will add two periodic jobs of integration test of helm with native k8s cluster, which use v2.2.0 chart-testing tool: 1. helm with 2.12.2 and kubernetes with v1.12.7 2. helm with 2.12.2 and kubernetes with v1.13.4 Closes: theopenlab/openlab#212 Closes: theopenlab/openlab#213
…e functions into projections
### What changes were proposed in this pull request?
This PR filters out `ExtractValues`s that contains any aggregation function in the `NestedColumnAliasing` rule to prevent cases where aggregations are pushed down into projections.
### Why are the changes needed?
To handle a corner/missed case in `NestedColumnAliasing` that can cause users to encounter a runtime exception.
Consider the following schema:
```
root
|-- a: struct (nullable = true)
| |-- c: struct (nullable = true)
| | |-- e: string (nullable = true)
| |-- d: integer (nullable = true)
|-- b: string (nullable = true)
```
and the query:
`SELECT MAX(a).c.e FROM (SELECT a, b FROM test_aggregates) GROUP BY b`
Executing the query before this PR will result in the error:
```
java.lang.UnsupportedOperationException: Cannot generate code for expression: max(input[0, struct<c:struct<e:string>,d:int>, true])
at org.apache.spark.sql.errors.QueryExecutionErrors$.cannotGenerateCodeForExpressionError(QueryExecutionErrors.scala:83)
at org.apache.spark.sql.catalyst.expressions.Unevaluable.doGenCode(Expression.scala:312)
at org.apache.spark.sql.catalyst.expressions.Unevaluable.doGenCode$(Expression.scala:311)
at org.apache.spark.sql.catalyst.expressions.aggregate.AggregateExpression.doGenCode(interfaces.scala:99)
...
```
The optimised plan before this PR is:
```
'Aggregate [b#1], [_extract_e#5 AS max(a).c.e#3]
+- 'Project [max(a#0).c.e AS _extract_e#5, b#1]
+- Relation default.test_aggregates[a#0,b#1] parquet
```
### Does this PR introduce _any_ user-facing change?
No
### How was this patch tested?
A new unit test in `NestedColumnAliasingSuite`. The test consists of the repro mentioned earlier.
The produced optimized plan is checked for equivalency with a plan of the form:
```
Aggregate [b#452], [max(a#451).c.e AS max('a)[c][e]#456]
+- LocalRelation <empty>, [a#451, b#452]
```
Closes #33921 from vicennial/spark-36677.
Authored-by: Venkata Sai Akhil Gudesa <[email protected]>
Signed-off-by: Liang-Chi Hsieh <[email protected]>
…e functions into projections
### What changes were proposed in this pull request?
This PR filters out `ExtractValues`s that contains any aggregation function in the `NestedColumnAliasing` rule to prevent cases where aggregations are pushed down into projections.
### Why are the changes needed?
To handle a corner/missed case in `NestedColumnAliasing` that can cause users to encounter a runtime exception.
Consider the following schema:
```
root
|-- a: struct (nullable = true)
| |-- c: struct (nullable = true)
| | |-- e: string (nullable = true)
| |-- d: integer (nullable = true)
|-- b: string (nullable = true)
```
and the query:
`SELECT MAX(a).c.e FROM (SELECT a, b FROM test_aggregates) GROUP BY b`
Executing the query before this PR will result in the error:
```
java.lang.UnsupportedOperationException: Cannot generate code for expression: max(input[0, struct<c:struct<e:string>,d:int>, true])
at org.apache.spark.sql.errors.QueryExecutionErrors$.cannotGenerateCodeForExpressionError(QueryExecutionErrors.scala:83)
at org.apache.spark.sql.catalyst.expressions.Unevaluable.doGenCode(Expression.scala:312)
at org.apache.spark.sql.catalyst.expressions.Unevaluable.doGenCode$(Expression.scala:311)
at org.apache.spark.sql.catalyst.expressions.aggregate.AggregateExpression.doGenCode(interfaces.scala:99)
...
```
The optimised plan before this PR is:
```
'Aggregate [b#1], [_extract_e#5 AS max(a).c.e#3]
+- 'Project [max(a#0).c.e AS _extract_e#5, b#1]
+- Relation default.test_aggregates[a#0,b#1] parquet
```
### Does this PR introduce _any_ user-facing change?
No
### How was this patch tested?
A new unit test in `NestedColumnAliasingSuite`. The test consists of the repro mentioned earlier.
The produced optimized plan is checked for equivalency with a plan of the form:
```
Aggregate [b#452], [max(a#451).c.e AS max('a)[c][e]#456]
+- LocalRelation <empty>, [a#451, b#452]
```
Closes #33921 from vicennial/spark-36677.
Authored-by: Venkata Sai Akhil Gudesa <[email protected]>
Signed-off-by: Liang-Chi Hsieh <[email protected]>
(cherry picked from commit 2ed6e7b)
Signed-off-by: Liang-Chi Hsieh <[email protected]>
… scale < 0 (apache#456)" (apache#467) This reverts commit 0017da5
…can reach ### What changes were proposed in this pull request? Task 225 (`PLAN_TASK_225.md`). For a change confined to Varka's files the module matrix is already off (task 160), so a fork run is four jobs: the sql Varka suites at 17 minutes, the catalyst suites at 11, and beside them Spark's "Linters, licenses, and dependencies" at 30 and "Documentation generation" at 16 (run 36295726808, apache#456). Every Varka PR waited on the lint job about a quarter of an hour past its own tests, and the two jobs held two of the fork's twenty slots for that long. **The lint job.** The precondition, which already classifies a change by its paths (`spark`, `scoped`, `docs`, `none`), now derives three flags, and the job's steps read them: - `lint-spark` is false for anything but a Spark change. It gates MiMa (8m08s: `MimaExcludes` excludes `org.apache.spark.sql.catalyst.*` and `org.apache.spark.sql.execution.*`, where Varka's Spark-side code lives, and the engine module has no released artifact to compare against), the dependency test (2m20s: it reads the poms, and a pom is a Spark file), the Connect client's MiMa (1m48s) and the R linter (3m: nothing of Varka's is under `R/`). - `lint-code` is false for documents. It gates the Scala and Java linters, which read Varka's sources. - `lint-python` is false unless a `.py` or `.pyi` file changed. It gates the Python linter (7m33s, most of it mypy over `python/pyspark`), which also lints Varka's tools under `dev/`. The license test, the JS linter, the config-policy check and the structured-logging check stay: under a minute together, and a new Varka file needs its header. A Spark change and the weekly full-matrix run keep every step, and a missing flag reads as true, so a caller that passes its own `jobs` gets the whole job. **The documentation job** is off for a `scoped` or `none` change, as it already was for `docs`: unidoc drops every source under `sql/catalyst`, `sql/execution` and `sql/internal` (`SparkBuild.ignoreUndocumentedPackages`) and leaves the engine project out of its filter, and the SQL docs list no internal config, so nothing Varka's code can change reaches the site. Any changed file under `docs/` keeps the job on whatever the verdict. That also closes a gap of task 160: `docs/sql-varka.md` counts as a Varka document, so a change to it alone skipped the job that renders it. **What it does not do:** change what runs for a Spark change or the weekly run; shorten the steps themselves; or trust the classifier beyond its rule (it reads paths, not content; a Varka-named file reaching into Spark's public API would pass MiMa here unchecked, and the review is the check for that, as for the matrix). ### Why are the changes needed? About 23 minutes of lint steps and a 16-minute documentation build per Varka PR that cannot fail on the change; with them gone, a scoped PR's wall time should be its own suites' 17 minutes, and two slots return to the pool sooner. ### Does this PR introduce _any_ user-facing change? No. CI only. ### How was this patch tested? - The workflow parses. The PR's own fork run is `spark` (it changes the workflow), so every step of both jobs runs in it: prediction 1 of `PLAN_TASK_225.md` 4. - The `scoped` and `docs` verdicts cannot be produced from this branch, since its diff carries the workflow change. The first scoped and documents-only pushes after the merge are the acceptance runs; predictions 2 to 4 give the expected times (about 7 minutes for the lint job on a scoped change without Python files, about 4 on documents, no documentation job), and section 5 records them. - `dev/varka_quote_check.py` (0 orphans) and the pre-commit hook. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Fable 5.1)
…g the CI slot ### What changes were proposed in this pull request? Task 226. On apache#456's first fork run, G14 of `VarkaCodegenGiveUpSuite` ran past 25 minutes, printing only ScalaTest's "still running", until the run was cancelled by hand. The cause (task 219's demotion loop) was fixed there; what this PR fixes is that nothing bounded it, so a regression of that kind holds the fork's single Build slot for the job's whole limit and names no test. **Why the existing cap did not act.** `SparkFunSuite` wraps every test in `failAfter(20 minutes)`, and it reads as a cap. It is not one: without a `Signaler`, ScalaTest's `failAfter` only checks the clock after the body returns, so a body that never returns is never reported. Interrupting the test thread would not end a CPU-bound loop either, since such a loop checks no interrupt, and `Thread.stop` is gone since JDK 20. **The watchdog.** `VarkaTestWatchdog`, a trait over `SparkFunSuite`, runs each test beside a watchdog thread. At the cap (`varka.test.watchdog.minutes`, ten by default) it interrupts the test thread, which ends a body that blocks or checks interrupts; thirty seconds later, if the test still runs, it writes every thread's stack under the test's name and halts the JVM. sbt then reports the suite aborted, and the log's last lines say which test hung and what it was doing, minutes after the hang instead of at the job's limit. The cap is sized from the record: the slowest Varka test, the width audit's census, takes about two minutes on the laptop, and the suites' CI jobs finish in eleven to seventeen minutes whole, so a single test past ten is a hang. A suite meant to run longer overrides `watchdogMinutes`. **Where it is mixed in.** Every Varka suite that extends `SparkFunSuite` or `QueryTest` directly (43), and the three Varka test bases (`VarkaEmitterTestBase`, `VarkaTpcCodegenCensus` for the TPC census suites). The two fuzzers are left out: their iteration counts bound them, and a nightly run of them is meant to take longer than ten minutes. The change to each suite is one `with VarkaTestWatchdog` and, outside the `varka` package, its import. **The test.** `VarkaTestWatchdogSuite` passes the mechanism its report sink and its halt as stubs, so the three outcomes take under a second each and the JVM stays: a body that returns is left alone, a body that blocks ends at the interrupt, and a body that spins ignoring the interrupt is reported with every thread's stack (the dump names the suite's own frame) and "halted". Records: row 226 marked done, and a lesson in `testing-and-debugging.md` with the `SKILLS.md` contents regenerated. ### Why are the changes needed? A hung test cost the fork a Build slot for the job's limit and a cancel by hand; now it costs ten minutes and leaves its name and its stacks in the log. ### Does this PR introduce _any_ user-facing change? No. Tests only. ### How was this patch tested? - `catalyst/testOnly *VarkaTestWatchdogSuite *VarkaLaneTypeSuite` and `sql/testOnly *VarkaInputRowsSuite *VarkaProjectExecSuite *VarkaTpcdsCodegenCensus*` pass (the suites that extend each base and each session provider with the watchdog mixed in); the fork CI runs every Varka suite with it. - The row's own check, reverting task 219's bisection on a branch and reading the failure, is what the fork CI would show on such a regression; the unit suite covers the mechanism without waiting ten minutes. - `dev/scalastyle` (which also sorted the added imports), the 100-column and non-ASCII scans, `dev/varka_quote_check.py` (0 orphans) and the pre-commit hook are clean. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Fable 5.1)
The goal is to improve the performance of the HiveTableScan Operator:
As a quick benchmark run the following code in the scala interpreter:
scala> :paste
hql("CREATE TABLE IF NOT EXISTS sample (key1 INT, key2 INT,value STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY ','")
hql("LOAD DATA LOCAL INPATH 'examples/src/main/resources/sample2.txt' INTO TABLE sample")
println("Result of SELECT * FROM sample:")
val start = System.nanoTime
val recs = hql("FROM sample SELECT key1,key2,value").collect()
val micros = (System.nanoTime - start) / 1000
println("%d microsecondss".format(micros))
scala> CTRL-D
you can download the test file from here:
http://homes.cs.washington.edu/~soroush/sample2.txt
"sample2.txt contains about 3.6 million rows. The improved code scans the entire table in about 9 seconds while the original code scans the entire table in about 22 seconds.
Regarding the last item in the task:
"Avoid Reading Unneeded Data - Some Hive Serializer/Deserializer (SerDe) interfaces support reading only the required columns from the underlying HDFS files. We should use ColumnProjectionUtils to configure these correctly."
The way to do it, should be similar to the following code:
https://github.com/amplab/shark/blob/master/src/main/scala/shark/execution/TableScanOperator.scala
I tried to take a similar approach, but I am not sure columnar reading is working at hiveOperators.scala right now. Anyway, it requires more time for me to make sure that last feature is working. Please notice that it was the first time that I wrote code in scala and it took me some time to get comfortable with the language.