【问题标题】:PySpark job seems to be stuck while materializing RDDPySpark 作业在实现 RDD 时似乎卡住了
【发布时间】:2016-06-30 04:00:24
【问题描述】:

我有一个 SparkJob,它首先在 N 个项目之间创建一个成对得分矩阵。虽然密集,但速度非常快,最多可达 20K 个元素,之后它似乎会卡住很长时间。我在多次尝试中看到的最后一条日志行是“已清理的累加器”,我附上了下面的代码块,以使用随机创建的 50K 元素数据集重现该问题。笛卡尔积非常快,对生成的 RDD 的计数会在几分钟内返回(25 亿行),但第二次计数会卡住两个多小时,日志或 Spark 作业 UI 中没有任何进度更新。我有一个由 15 个 EC2 M3.2xLarge 节点组成的集群。我如何才能了解这里发生了什么以及可以做些什么来加快速度?

import random
from pyspark.context import SparkContext
from pyspark.sql import HiveContext, SQLContext
import math
from pyspark.sql.types import *
from pyspark.sql.types import Row
sc=SparkContext(appName='kmedoids_test')
sqlContext=HiveContext(sc)
n=50000
A = [random.normalvariate(0, 1) for i in range(n)]
B = [random.normalvariate(1, 1) for i in range(n)]
C = [random.normalvariate(-1, 0.5) for i in range(n)]
df = sqlContext.createDataFrame(zip(A,B,C), ["A","B","C"])
f = lambda x, y : math.pow((x.A - y.A), 2) + math.pow((x.B - y.B), 2) +    math.pow((x.C - y.C), 2)
schema  = StructType([StructField("row_id", LongType(), False)] +       df.schema.fields[:])
no_of_cols=len(df.columns)
rdd_zipped_with_index=df.rdd.zipWithIndex()
reconstructed_rdd = rdd_zipped_with_index.map(lambda x: [x[1]]+list(x[0][0:no_of_cols]))
indexed_df=reconstructed_rdd.toDF(schema)
indexed_rdd = indexed_df.rdd
sc._conf.set("spark.sql.autoBroadcastJoinThreshold","-1") #turning off broadcast join
rdd_cartesian_prod = indexed_rdd.cartesian(indexed_rdd)
print "----------Count in self-join--------------- {0}".format(rdd_cartesian_prod.count()) #this returns quickly in about 160s
ScoreVec = Row("head_id","tail_id","score")
output_rdd = rdd_cartesian_prod.map(lambda x :   ScoreVec(float(x[0].row_id), float(x[1].row_id), float(f(x[0], x[1]))))
print "-----------Count after scoring---------------  {0}".format(output_rdd.count()) #gets stuck here for a LONG time
output_df = output_rdd.toDF() #does not get here

【问题讨论】:

  • "cleaned accumulator" 只是 Spark 不断吐出的一行,如果你不告诉它不那么冗长的话。

标签: apache-spark pyspark


【解决方案1】:

这可能是由于lazy evaluation

Spark 也是一样。它会一直等到你给它操作符,只有当你要求它给你最终答案时它才会评估,它总是看起来限制它必须做的工作量。

笛卡尔积之后的行数将是indexed_rdd.count()^2。 Spark 实际上并不需要生成所有这些行来知道会有多少行。尽管output_rdd.count() 中的行数相同,但Spark 实际上是在处理所有数据并在计数之前对其进行映射。这就是为什么这项任务需要更长的时间。要证明这是正在发生的事情,您可以尝试indexed_rdd.cache().count()。在计数之前缓存会强制数据处理(并将结果保存在内存中),并且会花费很长时间。

【讨论】:

  • 谢谢大卫。正如你所建议的,我尝试了 rdd_cartesian_prod.cache().count() 来强制处理自加入 RDD。仍然在一分钟内完成,而剩下的工作需要 45 分钟才能完成 30k 数据集。我怎样才能加快速度?增加并行性会适得其反并最终导致更多混乱吗?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多