Skip to content

[SPARK-21269][Core][WIP] Fix FetchFailedException when enable maxReqSizeShuffleToMem and KryoSerializer - #18490

Closed
wangyum wants to merge 1 commit into
apache:masterfrom
wangyum:missing-output-location
Closed

[SPARK-21269][Core][WIP] Fix FetchFailedException when enable maxReqSizeShuffleToMem and KryoSerializer#18490
wangyum wants to merge 1 commit into
apache:masterfrom
wangyum:missing-output-location

Conversation

@wangyum

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

Spark cluster can reproduce, local can't:

  1. Start a spark context with spark.reducer.maxReqSizeShuffleToMem=1K and spark.serializer=org.apache.spark.serializer.KryoSerializer:
$ spark-shell --conf spark.reducer.maxReqSizeShuffleToMem=1K --conf spark.serializer=org.apache.spark.serializer.KryoSerializer
  1. A shuffle:
scala> sc.parallelize(0 until 3000000, 10).repartition(2001).count()

The error messages:

org.apache.spark.shuffle.MetadataFetchFailedException: Missing an output location for shuffle 0
at org.apache.spark.MapOutputTracker$$anonfun$convertMapStatuses$2.apply(MapOutputTracker.scala:808)
at org.apache.spark.MapOutputTracker$$anonfun$convertMapStatuses$2.apply(MapOutputTracker.scala:804)
at scala.collection.TraversableLike$WithFilter$$anonfun$foreach$1.apply(TraversableLike.scala:733)
at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:186)
at scala.collection.TraversableLike$WithFilter.foreach(TraversableLike.scala:732)
at org.apache.spark.MapOutputTracker$.convertMapStatuses(MapOutputTracker.scala:804)
at org.apache.spark.MapOutputTrackerWorker.getMapSizesByExecutorId(MapOutputTracker.scala:618)
at org.apache.spark.shuffle.BlockStoreShuffleReader.read(BlockStoreShuffleReader.scala:49)
at org.apache.spark.rdd.ShuffledRDD.compute(ShuffledRDD.scala:105)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:323)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:287)
at org.apache.spark.rdd.CoalescedRDD$$anonfun$compute$1.apply(CoalescedRDD.scala:100)
at org.apache.spark.rdd.CoalescedRDD$$anonfun$compute$1.apply(CoalescedRDD.scala:99)
at scala.collection.Iterator$$anon$12.nextCur(Iterator.scala:434)
at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:440)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408)
at org.apache.spark.util.Utils$.getIteratorSize(Utils.scala:1802)
at org.apache.spark.rdd.RDD$$anonfun$count$1.apply(RDD.scala:1159)
at org.apache.spark.rdd.RDD$$anonfun$count$1.apply(RDD.scala:1159)
at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:2065)
at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:2065)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
at org.apache.spark.scheduler.Task.run(Task.scala:108)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:341)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
at java.lang.Thread.run(Thread.java:745)
)
...

This PR fix this issue.

How was this patch tested?

Add later

@SparkQA

Copy link
Copy Markdown

Test build #78987 has finished for PR 18490 at commit 1c837b6.

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

@zsxwing

zsxwing commented Jun 30, 2017

Copy link
Copy Markdown
Member

I don't get it. Could you point out which place serializes MapStatus using Kyro? NVM

@zsxwing

Copy link
Copy Markdown
Member

I cannot reproduce this issue. Could you provide a unit test to reproduce this?

Anyway, I suggest using kryo.register(classOf[HighlyCompressedMapStatus], new KryoJavaSerializer()) to force Kryo using Java serialization so that we don't need to worry about if the Kryo default serializer works or not.

@HyukjinKwonHyukjinKwon mentioned this pull request Jul 31, 2017
@asfgitasfgit closed this in 3a45c7fAug 5, 2017
@wangyum
wangyum deleted the missing-output-location branch September 11, 2017 22:44
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.

3 participants

@wangyum@SparkQA@zsxwing