【发布时间】:2019-11-07 01:55:54
【问题描述】:
您好,我正在使用 spark Mllib 并在 1M 数据集和 1k 数据集之间进行 approxSimilarityJoin。
当我这样做时,我会广播 1k 的。
我看到的是,工作在倒数第二个任务中停止。
所有的执行者都死了,但有一个会持续运行很长时间,直到内存不足。
我检查了神经节,它显示内存一直在上升,直到达到极限
磁盘空间一直在减少,直到完成:
我调用的操作是 write,但它对 count 的作用相同。
现在我想知道:是否有可能集群中的所有分区都汇聚到一个节点上并造成这个瓶颈?
这是我的代码 sn-p:
var dfW = cookesWb.withColumn("n", monotonically_increasing_id())
var bunchDf = dfW.filter(col("n").geq(0) && col("n").lt(1000000) )
bunchDf.repartition(3000)
model.
approxSimilarityJoin(bunchDf,broadcast(cookesNextLimited),80,"EuclideanDistance").
withColumn("min_distance", min(col("EuclideanDistance")).over(Window.partitionBy(col("datasetA.uid")))
).
filter(col("EuclideanDistance") === col("min_distance")).
select(col("datasetA.uid").alias("weboId"),
col("datasetB.nextploraId").alias("nextId"),
col("EuclideanDistance")).write.format("parquet").mode("overwrite").save("approxJoin.parquet")
【问题讨论】:
标签: apache-spark configuration pyspark cluster-computing