【问题标题】:Spark task hanging at [GC (Allocation Failure) ]Spark 任务挂在 [GC (Allocation Failure) ]
【发布时间】:2019-10-09 07:16:59
【问题描述】:

编辑: 注意:执行程序通常会发出消息 [GC (Allocation Failure) ] 。它运行它是因为它正在尝试为 Executor 分配内存,但是 executor 已满,因此它会在向 Executor 加载新内容时尝试 GC 以腾出空间。如果您的 Executor 在循环中执行此操作,则可能意味着您尝试加载到该 Executor 的内容太大。

我在 AWS EMR 5.8.0 上运行 Spark 2.2、Scala 2.11

我正在尝试对拒绝完成的数据集运行count 操作。令人沮丧的是,它只挂在一个特定的文件上。我在与 S3 不同的文件上运行此作业,没问题 - 它完全完成。原始 CSV 文件本身为 @18GB,我们对其运行转换以将原始 CSV 转换为 struct 列,为其提供一个额外的列。

我的环境的核心从属服务器是 8 个实例,每个实例是:

r3.2xlarge
16 vCore, 61 GiB memory, 160 SSD GB storage

我的 Spark 会话设置是:

implicit val spark = SparkSession
      .builder()
      .appName("MyApp")
      .master("yarn")
      .config("spark.speculation","false")
      .config("hive.metastore.uris", s"thrift://$hadoopIP:9083")
      .config("hive.exec.dynamic.partition", "true")
      .config("hive.exec.dynamic.partition.mode", "nonstrict")
      .config("mapreduce.fileoutputcommitter.algorithm.version", "2")
      .config("spark.dynamicAllocation.enabled", false)
      .config("spark.executor.cores", 5)
      .config("spark.executors.memory", "18G")
      .config("spark.yarn.executor.memoryOverhead", "2G")
      .config("spark.driver.memory", "18G")
      .config("spark.executor.instances", 23)
      .config("spark.default.parallelism", 230)
      .config("spark.sql.shuffle.partitions", 230)
      .enableHiveSupport()
      .getOrCreate()

数据来自 CSV 文件:

val ds = spark.read
          .option("header", "true")
          .option("delimiter", ",")
          .schema(/* 2 cols: [ValidatedNel, and a stuct schema */)
          .csv(sourceFromS3)
          .as(MyCaseClass)

val mappedDs:Dataset[ValidatedNel, MyCaseClass] = ds.map(...)

mappedDs.repartition(230)

val count = mappedDs.count() // never finishes

正如预期的那样,它启动了 230 个任务,完成了 229 个任务,除了中间某处的一个。见下文 - 第一个任务永远挂起,中间的任务完成没有问题(虽然很奇怪 - 大小记录/比率非常不同) - 其他 229 个任务看起来与完成的任务完全相同。

Index| ID |Attempt |Status|Locality Level|Executor ID / Host|                       Launch Time          |   Duration   |GC Time|Input Size / Records|Write Time | Shuffle Write Size / Records| Errors
110   117   0   RUNNING     RACK_LOCAL     11 / ip-XXX-XX-X-XX.uswest-2.compute.internal 2019/10/01 20:34:01    1.1 h   43 min     66.2 MB / 2289538                0.0 B / 0   
0     7     0   SUCCESS     PROCESS_LOCAL  9 / ip-XXX-XX-X-XXX.us-west-2.compute.internal 2019/10/01 20:32:10   1.0 s   16 ms      81.2 MB /293        5 ms         59.0 B / 1   <-- this task is odd, but finishes
1     8     0   SUCCESS     RACK_LOCAL      9 / ip-XXX-XX-X-XXX.us-west-2.compute.internal 2019/10/01 20:32:10  2.1 min     16 ms      81.2 MB /2894845        9 s          59.0 B / 1   <- the other tasks are all similar to this one

检查挂起任务的stdout,我反复看到以下永无止境:

2019-10-01T21:51:16.055+0000: [GC (Allocation Failure) 2019-10-01T21:51:16.055+0000: [ParNew: 10904K->0K(613440K), 0.0129982 secs]2019-10-01T21:51:16.068+0000: [CMS2019-10-01T21:51:16.099+0000: [CMS-concurrent-mark: 0.031/0.044 secs] [Times: user=0.17 sys=0.00, real=0.04 secs] 
 (concurrent mode failure): 4112635K->2940648K(4900940K), 0.4986233 secs] 4123539K->2940648K(5514380K), [Metaspace: 60372K->60372K(1103872K)], 0.5121869 secs] [Times: user=0.64 sys=0.00, real=0.51 secs] 

另一个注意事项是,在我调用计数之前,我先调用repartition(230),然后再调用Dataset[T] 上的count,以确保数据的平等分布

这是怎么回事?

【问题讨论】:

  • spark.executorS.memory=... - 是错字吗?我还会考虑缩减您的资源分配并摆脱所有不相关的配置选项。
  • spark.executors.memory= 是一个有效的设置 - 它是给每个执行程序多少 RAM。 docs.aws.amazon.com/emr/latest/ManagementGuide/… 。 5 个核心/执行器,每个实例 16 个核心 - 1 个用于守护进程 = 每个实例 15 个核心 / 5 = 每个实例 3 个执行器。实例的 61GiB 内存 / 3 个执行程序 = 20 -(61GiB 的 10% 用于开销)= 18。
  • @mazaneicha - 我需要所有这些设置,所以无法删除它们。您会建议配对哪些资源?
  • AFAIK 它的.executor.,单数,而不是复数。您要为 OS、YARN NodeManager、ResourceManager 以及节点上运行的其他任何东西留下 1 个核心和 1 GB“开销”?这对我来说听起来太激进了。
  • 我怀疑你需要 default.parallelism 和 shuffle.patitions 设置。不确定您的应用中还发生了什么,但如前所述,它也不需要推测或任何配置单元...选项。

标签: scala amazon-web-services apache-spark amazon-emr


【解决方案1】:

这可能与数据倾斜和/或数据解析问题有关。请注意,问题分区的记录数量远远多于已成功处理的记录:

Input Size /  Records
66.2 MB / 2289538
81.2 MB /293

我会检查所有分区文件的大小和记录数是否大致相同。问题或“好”分区文件中的行和/或列分隔符可能已关闭(对于 ~80 Mb 文件来说,293 行似乎太低了)。

【讨论】:

  • 好吧 - 在进行计数之前,数据被加载到 Dataset[T] 中,并且 T 字段都不是 Option[_],所以如果不是,解析就会失败。我再次查看了我的输入大小和记录-81.2MB/293 记录似乎已关闭,正如您所说-但这一步完成得很好。觉得很奇怪,因为其他 228 个任务都是@81.2MB/@2,200,000。我不会打电话给repartition(230) 正确地重新分配数据,防止出现偏差吗?
  • “加载”是什么意思? Spark 有一个惰性评估模型。所有读取/缓存/重新分区/数据转换调用仅用于构建执行计划 DAG。仅在遇到操作(如计数、写入、显示等)时执行。
  • 对不起,我的意思是在计数之前,有一个map 来自CSV带来的原始DF的转换,它被转换为一个Dataset[T],因为它是一个转换,它不会在那一步失败,而不是计数吗?
  • 我还应该提到,我在计数之前也为了调试目的做了一次 take(n),这也可以按预期工作。
  • 原来我试图对其执行操作的文件可能已损坏。通过 SSH 连接到 EMR 集群,我什至无法在这个东西上运行一个简单的 count。给你打勾,因为最终,这是一个数据解析问题!
猜你喜欢
  • 2015-05-15
  • 2021-04-28
  • 1970-01-01
  • 1970-01-01
  • 2020-09-30
  • 2021-08-08
  • 2018-03-19
  • 1970-01-01
  • 2015-08-19
相关资源
最近更新 更多