Repository navigation
[SPARK-1615] Synchronize accesses to the LiveListenerBus' event queue - #544
andrewor14 wants to merge 3 commits into
Conversation
This guards against the race condition in which we (1) dequeue an event, and (2) check for queue emptiness before (3) actually processing the event in all attached listeners. The solution is to make steps (1) and (3) atomic relatively to (2).
|
Merged build triggered. |
|
Merged build started. |
|
Merged build finished. |
|
Refer to this link for build results: https://amplab.cs.berkeley.edu/jenkins/job/SparkPullRequestBuilder/14462/ |
There was a problem hiding this comment.
Maybe we can only send release if eventAdded. It's a minor detail but would also mean that event should never be null, which should (slightly) help the bus to catch up with the workload.
There was a problem hiding this comment.
Yes, didn't realize offer also returns a boolean
|
LGTM, the semantics and performance of this seem quite reasonable. The synchronzied block should be virtually cost-less in normal operation, and I trust java.concurrent's implementation of Semaphore. |
|
Looks great. |
|
Merged build triggered. |
|
Merged build started. |
|
Merged build finished. All automated tests passed. |
|
All automated tests passed. |
|
Thanks - merged. |
Original poster is @zsxwing, who reported this bug in #516. Much of SparkListenerSuite relies on LiveListenerBus's `waitUntilEmpty()` method. As the name suggests, this waits until the event queue is empty. However, the following race condition could happen: (1) We dequeue an event (2) The queue is empty, we return true (even though the event has not been processed) (3) The test asserts something assuming that all listeners have finished executing (and fails) (4) The listeners receive and process the event This PR makes (1) and (4) atomic by synchronizing around it. To do that, however, we must avoid using `eventQueue.take`, which is blocking and will cause a deadlock if we synchronize around it. As a workaround, we use the non-blocking `eventQueue.poll` + a semaphore to provide the same semantics. This has been a possible race condition for a long time, but for some reason we've never run into it. Author: Andrew Or <andrewor14@gmail.com> Closes #544 from andrewor14/stage-info-test-fix and squashes the following commits: 3cbe40c [Andrew Or] Merge github.com:apache/spark into stage-info-test-fix 56dbbcb [Andrew Or] Check if event is actually added before releasing semaphore eb486ae [Andrew Or] Synchronize accesses to the LiveListenerBus' event queue (cherry picked from commit ee6f7e2) Signed-off-by: Patrick Wendell <pwendell@gmail.com>
…loses apache#544. Fixed warnings in test compilation. This commit fixes two problems: a redundant import, and a deprecated function. Author: Kay Ousterhout <kayousterhout@gmail.com> == Merge branch commits == commit da9d2e13ee4102bc58888df0559c65cb26232a82 Author: Kay Ousterhout <kayousterhout@gmail.com> Date: Wed Feb 5 11:41:51 2014 -0800 Fixed warnings in test compilation. This commit fixes two problems: a redundant import, and a deprecated function.
Original poster is @zsxwing, who reported this bug in apache#516. Much of SparkListenerSuite relies on LiveListenerBus's `waitUntilEmpty()` method. As the name suggests, this waits until the event queue is empty. However, the following race condition could happen: (1) We dequeue an event (2) The queue is empty, we return true (even though the event has not been processed) (3) The test asserts something assuming that all listeners have finished executing (and fails) (4) The listeners receive and process the event This PR makes (1) and (4) atomic by synchronizing around it. To do that, however, we must avoid using `eventQueue.take`, which is blocking and will cause a deadlock if we synchronize around it. As a workaround, we use the non-blocking `eventQueue.poll` + a semaphore to provide the same semantics. This has been a possible race condition for a long time, but for some reason we've never run into it. Author: Andrew Or <andrewor14@gmail.com> Closes apache#544 from andrewor14/stage-info-test-fix and squashes the following commits: 3cbe40c [Andrew Or] Merge github.com:apache/spark into stage-info-test-fix 56dbbcb [Andrew Or] Check if event is actually added before releasing semaphore eb486ae [Andrew Or] Synchronize accesses to the LiveListenerBus' event queue
Fix the stubbing of the reader benchmark tests
* Change log dir Change log dir * Update run.yaml * Update run.yaml * Update run.yaml
…API doesn't work as expected (apache#544)
…te` for Java 21 ### What changes were proposed in this pull request? SPARK-44507(#42130) updated `try_arithmetic.sql.out` and `numeric.sql.out`, SPARK-44868(#42534) updated `datetime-formatting.sql.out`, but these PRs didn’t pay attention to the test health on Java 21. So this PR has regenerated the golden files `try_arithmetic.sql.out.java21`, `numeric.sql.out.java21`, and `datetime-formatting.sql.out.java21` of `SQLQueryTestSuite` so that `SQLQueryTestSuite` can be tested with Java 21. ### Why are the changes needed? Restore `SQLQueryTestSuite` to be tested with Java 21. ### Does this PR introduce _any_ user-facing change? No ### How was this patch tested? - Pass GitHub Actions - Manual checked: ``` java -version openjdk version "21-ea" 2023-09-19 OpenJDK Runtime Environment Zulu21+69-CA (build 21-ea+28) OpenJDK 64-Bit Server VM Zulu21+69-CA (build 21-ea+28, mixed mode, sharing) ``` ``` SPARK_GENERATE_GOLDEN_FILES=0 build/sbt "sql/testOnly org.apache.spark.sql.SQLQueryTestSuite" ``` **Before** ``` ... [info] - datetime-formatting.sql *** FAILED *** (316 milliseconds) [info] datetime-formatting.sql [info] Array("-- Automatically generated by SQLQueryTestSuite [info] ", "create temporary view v as select col from values [info] (timestamp '1582-06-01 11:33:33.123UTC+080000'), [info] (timestamp '1970-01-01 00:00:00.000Europe/Paris'), [info] (timestamp '1970-12-31 23:59:59.999Asia/Srednekolymsk'), [info] (timestamp '1996-04-01 00:33:33.123Australia/Darwin'), [info] (timestamp '2018-11-17 13:33:33.123Z'), [info] (timestamp '2020-01-01 01:33:33.123Asia/Shanghai'), [info] (timestamp '2100-01-01 01:33:33.123America/Los_Angeles') t(col) [info] ", "struct<> [info] ", " [info] [info] [info] ", "select col, date_format(col, 'G GG GGG GGGG') from v [info] ", "struct<col:timestamp,date_format(col, G GG GGG GGGG):string> [info] ", "1582-05-31 19:40:35.123 AD AD AD Anno Domini [info] 1969-12-31 15:00:00 AD AD AD Anno Domini [info] 1970-12-31 04:59:59.999 AD AD AD Anno Domini [info] 1996-03-31 07:03:33.123 AD AD AD Anno Domini [info] 2018-11-17 05:33:33.123 AD AD AD Anno Domini [info] 2019-12-31 09:33:33.123 AD AD AD Anno Domini [info] 2100-01-01 01:33:33.123 AD AD AD Anno Domini [info] [info] [info] ", "select col, date_format(col, 'y yy yyy yyyy yyyyy yyyyyy') from v [info] ", "struct<col:timestamp,date_format(col, y yy yyy yyyy yyyyy yyyyyy):string> [info] ", "1582-05-31 19:40:35.123 1582 82 1582 1582 01582 001582 [info] 1969-12-31 15:00:00 1969 69 1969 1969 01969 001969 [info] 1970-12-31 04:59:59.999 1970 70 1970 1970 01970 001970 [info] 1996-03-31 07:03:33.123 1996 96 1996 1996 01996 001996 [info] 2018-11-17 05:33:33.123 2018 18 2018 2018 02018 002018 [info] 2019-12-31 09:33:33.123 2019 19 2019 2019 02019 002019 [info] 2100-01-01 01:33:33.123 2100 00 2100 2100 02100 002100 [info] ... [info] - postgreSQL/numeric.sql *** FAILED *** (35 seconds, 848 milliseconds) [info] postgreSQL/numeric.sql [info] Expected "...rg.apache.spark.sql.[]AnalysisException [info] { [info] ...", but got "...rg.apache.spark.sql.[catalyst.Extended]AnalysisException [info] { [info] ..." Result did not match for query #544 [info] SELECT '' AS to_number_2, to_number('-34,338,492.654,878', '99G999G999D999G999') (SQLQueryTestSuite.scala:876) [info] org.scalatest.exceptions.TestFailedException: ... [info] - try_arithmetic.sql *** FAILED *** (314 milliseconds) [info] try_arithmetic.sql [info] Expected "...rg.apache.spark.sql.[]AnalysisException [info] { [info] ...", but got "...rg.apache.spark.sql.[catalyst.Extended]AnalysisException [info] { [info] ..." Result did not match for query #20 [info] SELECT try_add(interval 2 year, interval 2 second) (SQLQueryTestSuite.scala:876) [info] org.scalatest.exceptions.TestFailedException: ``` **After** ``` [info] Run completed in 9 minutes, 10 seconds. [info] Total number of tests run: 572 [info] Suites: completed 1, aborted 0 [info] Tests: succeeded 572, failed 0, canceled 0, ignored 59, pending 0 [info] All tests passed. ``` ### Was this patch authored or co-authored using generative AI tooling? No Closes #42580 from LuciferYang/SPARK-44888. Authored-by: yangjie01 <yangjie01@baidu.com> Signed-off-by: yangjie01 <yangjie01@baidu.com>
…te` for Java 21 ### What changes were proposed in this pull request? SPARK-44507(apache#42130) updated `try_arithmetic.sql.out` and `numeric.sql.out`, SPARK-44868(apache#42534) updated `datetime-formatting.sql.out`, but these PRs didn’t pay attention to the test health on Java 21. So this PR has regenerated the golden files `try_arithmetic.sql.out.java21`, `numeric.sql.out.java21`, and `datetime-formatting.sql.out.java21` of `SQLQueryTestSuite` so that `SQLQueryTestSuite` can be tested with Java 21. ### Why are the changes needed? Restore `SQLQueryTestSuite` to be tested with Java 21. ### Does this PR introduce _any_ user-facing change? No ### How was this patch tested? - Pass GitHub Actions - Manual checked: ``` java -version openjdk version "21-ea" 2023-09-19 OpenJDK Runtime Environment Zulu21+69-CA (build 21-ea+28) OpenJDK 64-Bit Server VM Zulu21+69-CA (build 21-ea+28, mixed mode, sharing) ``` ``` SPARK_GENERATE_GOLDEN_FILES=0 build/sbt "sql/testOnly org.apache.spark.sql.SQLQueryTestSuite" ``` **Before** ``` ... [info] - datetime-formatting.sql *** FAILED *** (316 milliseconds) [info] datetime-formatting.sql [info] Array("-- Automatically generated by SQLQueryTestSuite [info] ", "create temporary view v as select col from values [info] (timestamp '1582-06-01 11:33:33.123UTC+080000'), [info] (timestamp '1970-01-01 00:00:00.000Europe/Paris'), [info] (timestamp '1970-12-31 23:59:59.999Asia/Srednekolymsk'), [info] (timestamp '1996-04-01 00:33:33.123Australia/Darwin'), [info] (timestamp '2018-11-17 13:33:33.123Z'), [info] (timestamp '2020-01-01 01:33:33.123Asia/Shanghai'), [info] (timestamp '2100-01-01 01:33:33.123America/Los_Angeles') t(col) [info] ", "struct<> [info] ", " [info] [info] [info] ", "select col, date_format(col, 'G GG GGG GGGG') from v [info] ", "struct<col:timestamp,date_format(col, G GG GGG GGGG):string> [info] ", "1582-05-31 19:40:35.123 AD AD AD Anno Domini [info] 1969-12-31 15:00:00 AD AD AD Anno Domini [info] 1970-12-31 04:59:59.999 AD AD AD Anno Domini [info] 1996-03-31 07:03:33.123 AD AD AD Anno Domini [info] 2018-11-17 05:33:33.123 AD AD AD Anno Domini [info] 2019-12-31 09:33:33.123 AD AD AD Anno Domini [info] 2100-01-01 01:33:33.123 AD AD AD Anno Domini [info] [info] [info] ", "select col, date_format(col, 'y yy yyy yyyy yyyyy yyyyyy') from v [info] ", "struct<col:timestamp,date_format(col, y yy yyy yyyy yyyyy yyyyyy):string> [info] ", "1582-05-31 19:40:35.123 1582 82 1582 1582 01582 001582 [info] 1969-12-31 15:00:00 1969 69 1969 1969 01969 001969 [info] 1970-12-31 04:59:59.999 1970 70 1970 1970 01970 001970 [info] 1996-03-31 07:03:33.123 1996 96 1996 1996 01996 001996 [info] 2018-11-17 05:33:33.123 2018 18 2018 2018 02018 002018 [info] 2019-12-31 09:33:33.123 2019 19 2019 2019 02019 002019 [info] 2100-01-01 01:33:33.123 2100 00 2100 2100 02100 002100 [info] ... [info] - postgreSQL/numeric.sql *** FAILED *** (35 seconds, 848 milliseconds) [info] postgreSQL/numeric.sql [info] Expected "...rg.apache.spark.sql.[]AnalysisException [info] { [info] ...", but got "...rg.apache.spark.sql.[catalyst.Extended]AnalysisException [info] { [info] ..." Result did not match for query apache#544 [info] SELECT '' AS to_number_2, to_number('-34,338,492.654,878', '99G999G999D999G999') (SQLQueryTestSuite.scala:876) [info] org.scalatest.exceptions.TestFailedException: ... [info] - try_arithmetic.sql *** FAILED *** (314 milliseconds) [info] try_arithmetic.sql [info] Expected "...rg.apache.spark.sql.[]AnalysisException [info] { [info] ...", but got "...rg.apache.spark.sql.[catalyst.Extended]AnalysisException [info] { [info] ..." Result did not match for query apache#20 [info] SELECT try_add(interval 2 year, interval 2 second) (SQLQueryTestSuite.scala:876) [info] org.scalatest.exceptions.TestFailedException: ``` **After** ``` [info] Run completed in 9 minutes, 10 seconds. [info] Total number of tests run: 572 [info] Suites: completed 1, aborted 0 [info] Tests: succeeded 572, failed 0, canceled 0, ignored 59, pending 0 [info] All tests passed. ``` ### Was this patch authored or co-authored using generative AI tooling? No Closes apache#42580 from LuciferYang/SPARK-44888. Authored-by: yangjie01 <yangjie01@baidu.com> Signed-off-by: yangjie01 <yangjie01@baidu.com>
### What changes were proposed in this pull request? Thirteen open rows of milestone 6 move to milestone 7, on the owner's decision of 2 October 2026 after the review of the milestone's open rows before its close. Each was read against the milestone's done-when list (`PLAN_MILESTONE_6.md` 1.3) and against the size-control rows, the goal set on 30 September, and none bears on either: - **180**, promotion, continuously: a cadence (the cross-posts, a findings post every two weeks, a living benchmark page) rather than a task the milestone could close; carried, as decided on 1 October. - **182**, extending Spark's own benchmarks: it was milestone 7's item 8, and returns to it. - **207 and 208**, two measured micro-optimisations of the range-set kernel. - **213**, the warm-up choosing which drivers to compile from statistics. - **214 to 217 and 224**, the rest of the Java port - the time, condition and chrono families, the facade, and the benchmark harness adapter with `VarkaEmitDump` - mechanical, and meant for an agent. - **218**, optional research into why the loop predicate cycles. - **222**, structural hashing of IR nodes, a compile-time cost. - **231**, the fallback path's row loops reading columns through the batch's vectors. The moves follow the precedent of 15 and 21 September (milestone 5's rows into items 15 and 39): every row keeps its text and number in `PLAN_MILESTONE_6.md` with a "Moved to milestone 7" note at its head, so every citation still resolves; sections 2.8 and 2.9, the design sections of 182 and 180, carry a note under their headings; section 8 lists the thirteen; and `SCOPE_MILESTONE_7.md` item 81 tables them with where each design is and why each left. Row 181, the second post, published on 1 October, is marked done, the way row 210 was for the first. What is left open in milestone 6 after this and the pull requests in flight (apache#544 closing 233, apache#546 closing 235, apache#547 closing 223 and adding 239): the two size-control rows, 236 (plan a kernel's size before building it) and 239 (locals a loop method stores and never reads). ### Why are the changes needed? So that milestone 6 can close on what it set out to do, and the moved work keeps its place and its reasons. ### Does this PR introduce _any_ user-facing change? No. Two planning documents. ### How was this patch tested? `dev/varka_quote_check.py` reports zero orphans and `dev/varka_precommit.sh` no findings. ### Was this patch authored or co-authored using generative AI tooling? Generated-by: Claude Code (Claude Opus 5.5)
Brings apache#544 (task 233's post) and apache#550 (thirteen rows moved to milestone 7); no conflicts. Generated-by: Claude Code (Claude Opus 5.5)
Brings apache#544 (task 233's post) and apache#550 (thirteen rows moved to milestone 7). PLAN_MILESTONE_6.md conflicted where the moved rows 222 and 224 sit beside row 223, and where row 239 sits beside row 181: each row takes its own side, 222, 224 and 181 master's, 223 and 239 this branch's. Generated-by: Claude Code (Claude Opus 5.5)
Brings apache#544, on which this branch was built, and apache#550 (thirteen rows moved to milestone 7). PLAN_TASK_233.md and POST_MILESTONE_6_GIVEUPS.md conflicted only where this branch had already changed apache#544's lines - section 17 and the two corrected sentences - so this branch's text is kept; PLAN_MILESTONE_6.md is master's, with this branch's note on row 233. Generated-by: Claude Code (Claude Opus 5.5)
### What changes were proposed in this pull request?
Task 233's four laptop benchmarks, the ones post 3 ("When Spark stops compiling your query") quotes beside the runners, ran unpinned: plain `build/sbt runMain`, so the benchmark thread could move between the laptop's four Zen 5 cores at 5.16 GHz and its eight Zen 5c cores at 3.29 GHz. In the owner's idle hour they ran again on the quiet laptop, pinned to the fast cores the way `dev/varka_bench_regen.sh` pins them (`taskset -c 0-3,12-15`), and were repeated where the post quotes them: the fallback-cost and interpreter benchmarks five times, the wide projection twice (its first pinned run moved under 3% on every row), and the compile wait once. Each `*-jdk25-laptop-pinned-results.txt` holds every run as the benchmark printed it, with its provenance, including that on eight processors the JVM starts 4 C2 compiler threads instead of 12, and which runs had other work on the other core complex.
`PLAN_TASK_233.md` 17 scores every laptop number the post quotes against them, by a rule set before the repeats: a published number is corrected only where all five pinned runs fall outside it.
- The wide-projection table stands: 1.16 and 1.25 times slower for 50 and 99 cheap columns against 1.18 to 1.20 and 1.24 to 1.27 pinned, 12% and 16% faster for mixed columns against 11 to 17% and 16 to 18%. So does the time C2's code arrives in Figure 3 (13 s; 13.0 to 14.4 s pinned).
- The 150-column stage after C2 is corrected. The post said it stays 1.06 times the time outside a stage on the laptop; the five pinned runs give 0.74 to 1.00. The stage's own time is steady (406 to 423 ms); the out-of-stage time is what moves, 407 to 564 ms between runs and, within run 1, 385 in its table against 564 in its series, a JIT outcome per JVM rather than the pin. The post now says "on the laptop, 0.7 to 1.1 times, depending on the run", which covers all six runs, and Figure 3's laptop line, one of them, stays.
- The interpreter range is corrected: the post said 3.6 to 4.4 times the compiled projection on the laptop; the pinned runs give 3.06 to 3.44 at 100 columns and 3.61 to 4.20 at 1000, below the unpinned run at every width. The post now says 3.1 to 4.2.
Row 233 notes the correction. **Built on apache#544; merge apache#544 first.** **Republished** on the owner's go: the page rebuilt from this branch is live as `708fd21` in `vecbricks/vecbricks.github.io`, and the last commit here records it.
### Why are the changes needed?
A laptop number in a published post should not depend on which cores the scheduler happened to pick. The rerun was deferred to the owner's next idle-machine window, which came today.
### Does this PR introduce _any_ user-facing change?
No. Results files, the plan, and two sentences of a post.
### How was this patch tested?
Each benchmark ran one at a time with the load average under 1.0 at its start (recorded per run). `dev/varka_quote_check.py` reports zero orphans; `dev/varka_precommit.sh` reports no findings. The page rendered from this branch with `dev/varka_post_page.py` differs from the live one only in the two sentences.
### Was this patch authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus 5.5)
Original poster is @zsxwing, who reported this bug in #516.
Much of SparkListenerSuite relies on LiveListenerBus's
waitUntilEmpty()method. As the name suggests, this waits until the event queue is empty. However, the following race condition could happen:(1) We dequeue an event
(2) The queue is empty, we return true (even though the event has not been processed)
(3) The test asserts something assuming that all listeners have finished executing (and fails)
(4) The listeners receive and process the event
This PR makes (1) and (4) atomic by synchronizing around it. To do that, however, we must avoid using
eventQueue.take, which is blocking and will cause a deadlock if we synchronize around it. As a workaround, we use the non-blockingeventQueue.poll+ a semaphore to provide the same semantics.This has been a possible race condition for a long time, but for some reason we've never run into it.