Skip to content

[SPARK-21099][Spark Core] INFO Log Message Using Incorrect Executor I… - #18308

Closed
ihazem wants to merge 3 commits into
apache:masterfrom
ihazem:master
Closed

[SPARK-21099][Spark Core] INFO Log Message Using Incorrect Executor I…#18308
ihazem wants to merge 3 commits into
apache:masterfrom
ihazem:master

Conversation

@ihazem

Copy link
Copy Markdown

…dle Timeout Value

What changes were proposed in this pull request?

Updated the INFO logging method (logInfo) to use $timeout instead, which uses cachedExecutorIdleTimeoutS if the executor has cached blocks

How was this patch tested?

Testing was done by doing the following:

  1. Update spark-defaults.conf to set the following:
    executorIdleTimeout=30
    cachedExecutorIdleTimeout=20
  2. Update log4j.properties to set the following:
    shell.log.level=INFO
  3. Run the following in spark-shell:
    scala> val textFile = sc.textFile("/user/spark/applicationHistory/app_1234")
    scala> textFile.cache().count()
  4. After 20 secs I see the timeout message for idle cache, and after 30 secs I see the timeout message for executor (executorIdleTimeoutS).

Please review http://spark.apache.org/contributing.html before opening a pull request.

@jerryshao

Copy link
Copy Markdown
Contributor

LGTM. BTW can you please complement the PR description, thanks.

@srowen

Copy link
Copy Markdown
Member

Jenkins test this please

@SparkQA

Copy link
Copy Markdown

Test build #78086 has finished for PR 18308 at commit 00a42e7.

  • This patch fails to build.
  • This patch merges cleanly.
  • This patch adds no public classes.

if (testing || executorsRemoved.nonEmpty) {
executorsRemoved.foreach { removedExecutorId =>
newExecutorTotal -= 1
val hasCachedBlocks = SparkEnv.get.blockManager.master.hasCachedBlocks(executorId);

@jerryshaojerryshaoJun 15, 2017

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.

This variable executorId is not defined, should change to removedExecutorId.

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.

And final semicolon ";" is not necessary.

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.

Btw, there could be a chance when querying executor from BlockManager, the executor/block manager was already removed, so we will potentially get false.

@ihazem

Copy link
Copy Markdown
Author

@jerryshao - I pushed another commit fixing removedExecutorId and removing ';'...thx for feedback on that.

Regarding the edge case where the block manager was already removed and could potentially get false. then this:
 val timeout = if (hasCachedBlocks) cachedExecutorIdleTimeoutS else executorIdleTimeoutS
would just evaluate to executorIdleTimeoutS and wouldn't matter anyway since executor is gone. I wonder if whether executor is completely gone or whether executor is still there but has no cached RDD, if both scenarios return false. That seems to be the case (unless I'm reading it wrong) based on HasCachedBlocks in core/src/main/scala/org/apache/spark/storage/BlockManagerMasterEndpoint.scala

@vanzin

Copy link
Copy Markdown
Contributor

ok to test

@SparkQA

Copy link
Copy Markdown

Test build #78126 has finished for PR 18308 at commit 0f1c467.

  • This patch fails Spark unit tests.
  • This patch merges cleanly.
  • This patch adds no public classes.

@jerryshao

Copy link
Copy Markdown
Contributor

I wonder if whether executor is completely gone or whether executor is still there but has no cached RDD, if both scenarios return false.

Yes, that's the case, we cannot differentiate this two scenarios. But I think it is fine, since it is just a log issue and hard for us to differentiate them in the current code.

@SparkQA

Copy link
Copy Markdown

Test build #3799 has finished for PR 18308 at commit 0f1c467.

  • This patch passes all tests.
  • This patch merges cleanly.
  • This patch adds no public classes.

if (testing || executorsRemoved.nonEmpty) {
executorsRemoved.foreach { removedExecutorId =>
newExecutorTotal -= 1
val hasCachedBlocks = SparkEnv.get.blockManager.master.hasCachedBlocks(removedExecutorId)

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.

This is a pretty expensive call just so you can print a more accurate log message... at the very least it should be inside the logInfo call:

logInfo {
expensiveStuff()
"This is the log message referencing the expensive stuff."
}

But I'd prefer a solution that doesn't make this call at all, since INFO is more often than not enabled.

@AmplabJenkins

Copy link
Copy Markdown

Can one of the admins verify this patch?

@srowen

Copy link
Copy Markdown
Member

Ping @ihazem

@srowen

Copy link
Copy Markdown
Member

Ping @ihazem to update or close

@ihazem

Copy link
Copy Markdown
Author

@srowen - I modified to take @vanzin recommendation of putting the logic in the logInfo call instead. I'm open to hearing any recommendations on using a less expensive call. Thx!

val timeout = if (hasCachedBlocks) cachedExecutorIdleTimeoutS else executorIdleTimeoutS
logInfo(s"Removing executor $removedExecutorId because it has been idle for " +
s"$executorIdleTimeoutS seconds (new desired total will be $newExecutorTotal)")
s"${if (SparkEnv.get.blockManager.master.hasCachedBlocks(executorId)) cachedExecutorIdleTimeoutS else executorIdleTimeoutS} seconds (new desired total will be $newExecutorTotal)")

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.

This is hard to read. @vanzin 's idea is better. Can you split this into two statements? cheap one at info, expensive one at debug?

@jiangxb1987

Copy link
Copy Markdown
Contributor

Can you update the title to:

[SPARK-21099][Core] Log cachedExecutorIdleTimeoutS instead of executorIdleTimeoutS if the executor has cached blocks

?

@jerryshao

jerryshao commented Jul 3, 2017

Copy link
Copy Markdown
Contributor

Is this the final modified code @ihazem , why do you have hasCachedBlocks logics both inside and outside of logInfo statement? Also the code is too long.

Can you please at least do a round of self-review before pushing the changes.

@ihazem

Copy link
Copy Markdown
Author

@jerryshao - sorry about that...let me remove the logic outside the logInfo statements. Missed that. I'll do some more self-review and peer-review going fwd before submitting to not waste everyone's time here.

@jiangxb1987 - i'll update title per your recommendation

@srowen - do you or @vanzin have any recommendations on making a cheaper call in the info level? I can definitely make the more expensive call at debug/warn, but could use some direction on a cheaper call for the info level. Basically, we want to use the cachedExecutorIdleTimeoutS if the executor has cached blocks (whereas right now it uses executorIdleTimeoutS).

Thx!

@vanzin

Copy link
Copy Markdown
Contributor

There is no cheaper call you can make. Which is the problem here. So either you have to modify the code so that it keeps track of why the executor has timed out (i.e. more state being passed around), or you need to tweak the message so that it says that the timeout is different when there are cached blocks on the executor.

I'm not so sure there's a lot to gain from printing the exact timeout value in this message, so I'm leaning towards just tweaking the message a little bit (e.g. add "or $otherTimeout ms if executor had cached blocks"), or even not doing anything.

@srowensrowen mentioned this pull request Jul 31, 2017
@asfgitasfgit closed this in 3a45c7fAug 5, 2017
zifeif2 pushed a commit to zifeif2/spark that referenced this pull request Nov 22, 2025
## What changes were proposed in this pull request?
This PR proposes to close stale PRs, mostly the same instances with apache#18017Closesapache#14085 - [SPARK-16408][SQL] SparkSQL Added file get Exception: is a directory …
Closesapache#14239 - [SPARK-16593] [CORE] [WIP] Provide a pre-fetch mechanism to accelerate shuffle stage.
Closesapache#14567 - [SPARK-16992][PYSPARK] Python Pep8 formatting and import reorganisation
Closesapache#14579 - [SPARK-16921][PYSPARK] RDD/DataFrame persist()/cache() should return Python context managers
Closesapache#14601 - [SPARK-13979][Core] Killed executor is re spawned without AWS key…
Closesapache#14830 - [SPARK-16992][PYSPARK][DOCS] import sort and autopep8 on Pyspark examples
Closesapache#14963 - [SPARK-16992][PYSPARK] Virtualenv for Pylint and pep8 in lint-python
Closesapache#15227 - [SPARK-17655][SQL]Remove unused variables declarations and definations in a WholeStageCodeGened stage
Closesapache#15240 - [SPARK-17556] [CORE] [SQL] Executor side broadcast for broadcast joins
Closesapache#15405 - [SPARK-15917][CORE] Added support for number of executors in Standalone [WIP]
Closesapache#16099 - [SPARK-18665][SQL] set statement state to "ERROR" after user cancel job
Closesapache#16445 - [SPARK-19043][SQL]Make SparkSQLSessionManager more configurable
Closesapache#16618 - [SPARK-14409][ML][WIP] Add RankingEvaluator
Closesapache#16766 - [SPARK-19426][SQL] Custom coalesce for Dataset
Closesapache#16832 - [SPARK-19490][SQL] ignore case sensitivity when filtering hive partition columns
Closesapache#17052 - [SPARK-19690][SS] Join a streaming DataFrame with a batch DataFrame which has an aggregation may not work
Closesapache#17267 - [SPARK-19926][PYSPARK] Make pyspark exception more user-friendly
Closesapache#17371 - [SPARK-19903][PYSPARK][SS] window operator miss the `watermark` metadata of time column
Closesapache#17401 - [SPARK-18364][YARN] Expose metrics for YarnShuffleService
Closesapache#17519 - [SPARK-15352][Doc] follow-up: add configuration docs for topology-aware block replication
Closesapache#17530 - [SPARK-5158] Access kerberized HDFS from Spark standalone
Closesapache#17854 - [SPARK-20564][Deploy] Reduce massive executor failures when executor count is large (>2000)
Closesapache#17979 - [SPARK-19320][MESOS][WIP]allow specifying a hard limit on number of gpus required in each spark executor when running on mesos
Closesapache#18127 - [SPARK-6628][SQL][Branch-2.1] Fix ClassCastException when executing sql statement 'insert into' on hbase table
Closesapache#18236 - [SPARK-21015] Check field name is not null and empty in GenericRowWit…
Closesapache#18269 - [SPARK-21056][SQL] Use at most one spark job to list files in InMemoryFileIndex
Closesapache#18328 - [SPARK-21121][SQL] Support changing storage level via the spark.sql.inMemoryColumnarStorage.level variable
Closesapache#18354 - [SPARK-18016][SQL][CATALYST][BRANCH-2.1] Code Generation: Constant Pool Limit - Class Splitting
Closesapache#18383 - [SPARK-21167][SS] Set kafka clientId while fetch messages
Closesapache#18414 - [SPARK-21169] [core] Make sure to update application status to RUNNING if executors are accepted and RUNNING after recovery
Closesapache#18432 - resolve com.esotericsoftware.kryo.KryoException
Closesapache#18490 - [SPARK-21269][Core][WIP] Fix FetchFailedException when enable maxReqSizeShuffleToMem and KryoSerializer
Closesapache#18585 - SPARK-21359
Closesapache#18609 - Spark SQL merge small files to big files Update InsertIntoHiveTable.scala
Added:
Closesapache#18308 - [SPARK-21099][Spark Core] INFO Log Message Using Incorrect Executor I…
Closesapache#18599 - [SPARK-21372] spark writes one log file even I set the number of spark_rotate_log to 0
Closesapache#18619 - [SPARK-21397][BUILD]Maven shade plugin adding dependency-reduced-pom.xml to …
Closesapache#18667 - Fix the simpleString used in error messages
Closesapache#18782 - Branch 2.1
Added:
Closesapache#17694 - [SPARK-12717][PYSPARK] Resolving race condition with pyspark broadcasts when using multiple threads
Added:
Closesapache#16456 - [SPARK-18994] clean up the local directories for application in future by annother thread
Closesapache#18683 - [SPARK-21474][CORE] Make number of parallel fetches from a reducer configurable
Closesapache#18690 - [SPARK-21334][CORE] Add metrics reporting service to External Shuffle Server
Added:
Closesapache#18827 - Merge pull request 1 from apache/master
## How was this patch tested?
N/A
Author: hyukjinkwon <gurwls223@gmail.com>
Closesapache#18780 from HyukjinKwon/close-prs.
Sign up for freeto 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.

7 participants

@ihazem@jerryshao@srowen@SparkQA@vanzin@AmplabJenkins@jiangxb1987