【发布时间】: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