Hri*_*lov 5 java amazon-web-services apache-spark
我的 Spark 应用程序在内存使用方面遇到了大问题。因此,我在集群模式下的 aws emr 集群上运行 Spark(这意味着它使用纱线作为集群管理器)。我无法提供应用程序的任何代码,因为它很大,我不知道问题出在哪里,但基本上我所做的是从 aws s3 加载几种类型的 json 对象(数据按小时存储在不同的文件夹中,即 2017/05/31/00 是 Spark 应用程序应在 2017 年 5 月 31 日 00:00 处理的数据,我每小时运行该应用程序),将它们反序列化为 Java pojo 对象,制作一些转换并缓存它们,因为我不想在每个操作上重新加载它们。一些对象被转换为java映射,并且在所有逻辑完成之后,java映射被转换回RDD并保存回s3中。
问题一:
具有以下配置:
-conf Spark.yarn.executor.memoryOverhead=3000 --driver-内存 8G --executor-内存 8G --deploy-mode 集群
我转换为 RDD 并保存它的 Map 的大小随着执行的每一小时而增加。但在完全随机的某些执行中,保存的输出丢失了数据。例如,如果成功执行后 json 对象的计数为 1000,而下一次执行时可能会变为 900,而不是增加。如果您在同一小时再次运行它,那么绝对有可能没问题(我在我的应用程序中看不到错误,因为这种情况很少发生并且是随机的)。所以这让我认为这是一个内存问题,也许内存中的某些对象被垃圾收集或类似的东西。
问题2:
当我使用以下配置运行程序时:
-conf Spark.yarn.executor.memoryOverhead=3000 --driver-内存 8G --executor-内存1G --deploy-mode 集群
在一次执行中几个小时(在应用程序中每小时执行一次 while 循环)然后一切都很好。但是,如果我在第三个小时将执行程序内存更改为超过 1G,则会抛出一个非常奇怪的异常:
17/05/31 09:38:52 INFO S3NativeFileSystem: Encountered an exception while reading 'data/appletv/user_mapping/part-00000.gz', will retry by attempting to reopen stream.
java.io.EOFException: Invalid position: 1856, exceeds the bounds of the stream: [0, 1855]
at com.amazon.ws.emr.hadoop.fs.s3n.S3NativeFileSystem$NativeS3FsInputStream.throwPositionOutOfBoundsException(S3NativeFileSystem.java:254)
at com.amazon.ws.emr.hadoop.fs.s3n.S3NativeFileSystem$NativeS3FsInputStream.retrievePair(S3NativeFileSystem.java:238)
at com.amazon.ws.emr.hadoop.fs.s3n.S3NativeFileSystem$NativeS3FsInputStream.reopenStream(S3NativeFileSystem.java:291)
at com.amazon.ws.emr.hadoop.fs.s3n.S3NativeFileSystem$NativeS3FsInputStream.read(S3NativeFileSystem.java:187)
at java.io.BufferedInputStream.read1(BufferedInputStream.java:284)
at java.io.BufferedInputStream.read(BufferedInputStream.java:345)
at java.io.DataInputStream.read(DataInputStream.java:149)
at org.apache.hadoop.io.compress.DecompressorStream.getCompressedData(DecompressorStream.java:159)
at org.apache.hadoop.io.compress.DecompressorStream.decompress(DecompressorStream.java:143)
at org.apache.hadoop.io.compress.DecompressorStream.read(DecompressorStream.java:85)
at java.io.InputStream.read(InputStream.java:101)
at org.apache.hadoop.util.LineReader.fillBuffer(LineReader.java:180)
at org.apache.hadoop.util.LineReader.readDefaultLine(LineReader.java:216)
at org.apache.hadoop.util.LineReader.readLine(LineReader.java:174)
at org.apache.hadoop.mapreduce.lib.input.LineRecordReader.nextKeyValue(LineRecordReader.java:186)
at org.apache.spark.rdd.NewHadoopRDD$$anon$1.hasNext(NewHadoopRDD.scala:199)
at org.apache.spark.InterruptibleIterator.hasNext(InterruptibleIterator.scala:39)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408)
at scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:408)
at org.apache.spark.shuffle.sort.BypassMergeSortShuffleWriter.write(BypassMergeSortShuffleWriter.java:149)
at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:96)
at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:53)
at org.apache.spark.scheduler.Task.run(Task.scala:99)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:282)
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)
Run Code Online (Sandbox Code Playgroud)
我实在找不到这两个问题的逻辑解释。如果有人能提供帮助那就太好了。