Last few tasks of Spark hang and progress extremely slowly

Viewed 522

My spark job reads from and writes to a S3 bucket, the problem I had lately is that last few tasks of a job just hang there, no error thrown, just progress very very slowly. enter image description here

enter image description here

e.g. these last 7 tasks take increasingly long time to finish, you can see Median is 28s, but Max for last few tasks can be high as 40 mins. I'm pretty sure data is not skewed, each partition contains only ~100Mb data, but still takes 40 mins to finish as shown in below screenshot. enter image description here

My question is what could be potential causes for this problem? Could this because of S3 connection etc., because we used to run spark jobs on NFS and didn't have such issue before. I have tried almost all configurations I can possibly find to help resolve this issue, but still no luck :(

This is a copy of thread dump of one of those "hanging" executors:

java.base@11.0.10/java.net.SocketInputStream.socketRead0(Native Method)
java.base@11.0.10/java.net.SocketInputStream.socketRead(Unknown Source)
java.base@11.0.10/java.net.SocketInputStream.read(Unknown Source)
java.base@11.0.10/java.net.SocketInputStream.read(Unknown Source)
app//com.amazonaws.thirdparty.apache.http.impl.io.SessionInputBufferImpl.streamRead(SessionInputBufferImpl.java:137)
app//com.amazonaws.thirdparty.apache.http.impl.io.SessionInputBufferImpl.read(SessionInputBufferImpl.java:197)
app//com.amazonaws.thirdparty.apache.http.impl.io.ContentLengthInputStream.read(ContentLengthInputStream.java:176)
app//com.amazonaws.thirdparty.apache.http.conn.EofSensorInputStream.read(EofSensorInputStream.java:135)
app//com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
app//com.amazonaws.event.ProgressInputStream.read(ProgressInputStream.java:180)
app//com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
app//com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
app//com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
app//com.amazonaws.event.ProgressInputStream.read(ProgressInputStream.java:180)
app//com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
app//com.amazonaws.util.LengthCheckInputStream.read(LengthCheckInputStream.java:107)
app//com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
app//com.amazonaws.services.s3.internal.S3AbortableInputStream.read(S3AbortableInputStream.java:125)
app//com.amazonaws.internal.SdkFilterInputStream.read(SdkFilterInputStream.java:90)
app//org.apache.hadoop.fs.s3a.S3AInputStream.lambda$read$3(S3AInputStream.java:468)
app//org.apache.hadoop.fs.s3a.S3AInputStream$$Lambda$1124/0x00007f862e16fd08.execute(Unknown Source)
app//org.apache.hadoop.fs.s3a.Invoker.once(Invoker.java:109)
app//org.apache.hadoop.fs.s3a.Invoker.lambda$retry$3(Invoker.java:265)
app//org.apache.hadoop.fs.s3a.Invoker$$Lambda$954/0x00007f86647c4900.execute(Unknown Source)
app//org.apache.hadoop.fs.s3a.Invoker.retryUntranslated(Invoker.java:322)
app//org.apache.hadoop.fs.s3a.Invoker.retry(Invoker.java:261)
app//org.apache.hadoop.fs.s3a.Invoker.retry(Invoker.java:236)
app//org.apache.hadoop.fs.s3a.S3AInputStream.read(S3AInputStream.java:464) => holding Monitor(org.apache.hadoop.fs.s3a.S3AInputStream@1670038891})
java.base@11.0.10/java.io.DataInputStream.read(Unknown Source)
app//org.apache.parquet.io.DelegatingSeekableInputStream.readFully(DelegatingSeekableInputStream.java:102)
app//org.apache.parquet.io.DelegatingSeekableInputStream.readFullyHeapBuffer(DelegatingSeekableInputStream.java:127)
app//org.apache.parquet.io.DelegatingSeekableInputStream.readFully(DelegatingSeekableInputStream.java:91)
app//org.apache.parquet.hadoop.ParquetFileReader$ConsecutiveChunkList.readAll(ParquetFileReader.java:1174)
app//org.apache.parquet.hadoop.ParquetFileReader.readNextRowGroup(ParquetFileReader.java:805)
app//org.apache.spark.sql.execution.datasources.parquet.VectorizedParquetRecordReader.checkEndOfRowGroup(VectorizedParquetRecordReader.java:313)
app//org.apache.spark.sql.execution.datasources.parquet.VectorizedParquetRecordReader.nextBatch(VectorizedParquetRecordReader.java:268)
app//org.apache.spark.sql.execution.datasources.parquet.VectorizedParquetRecordReader.nextKeyValue(VectorizedParquetRecordReader.java:174)
app//org.apache.spark.sql.execution.datasources.RecordReaderIterator.hasNext(RecordReaderIterator.scala:39)
app//org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.hasNext(FileScanRDD.scala:93)
app//org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.nextIterator(FileScanRDD.scala:173)
app//org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.hasNext(FileScanRDD.scala:93)
app//scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:458)
app//scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:458)
app//scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:458)
app//org.apache.spark.shuffle.sort.UnsafeShuffleWriter.write(UnsafeShuffleWriter.java:177)
app//org.apache.spark.shuffle.ShuffleWriteProcessor.write(ShuffleWriteProcessor.scala:59)
app//org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99)
app//org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:52)
app//org.apache.spark.scheduler.Task.run(Task.scala:127)
app//org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:446)
app//org.apache.spark.executor.Executor$TaskRunner$$Lambda$650/0x00007f86653e2440.apply(Unknown Source)
app//org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1377)
app//org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:449)
java.base@11.0.10/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)
java.base@11.0.10/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)
java.base@11.0.10/java.lang.Thread.run(Unknown Source)
0 Answers
Related