【发布时间】:2022-08-19 04:00:37
【问题描述】:
我有一个包含 1 亿行和 10,000 列的数据框。这些列有两种类型,标准 (C_i) 和动态 (X_i)。这个dataframe是经过一些处理得到的,性能很快。现在只剩下两个步骤:
目标:
- 需要使用相同的 C_i 列子集对每个 X_i 执行特定操作。
- 将每个 X-i 列转换为
FloatType。困难:
- 随着列数的增加,性能会严重下降。
- 一段时间后,似乎只有 1 个执行程序可以工作(%CPU 使用率 < 200%),即使是在具有 100 行和 1,000 列的样本数据上也是如此。如果我将它推到 1,500 列,它就会崩溃。
最小代码:
import spark.implicits._ import org.apache.spark.sql.types.FloatType // sample_udf val foo = (s_val: String, t_val: String) => { t_val + s_val.takeRight(1) } val foos_udf = udf(foo) spark.udf.register(\"foos_udf\", foo) val columns = Seq(\"C1\", \"C2\", \"X1\", \"X2\", \"X3\", \"X4\") val data = Seq((\"abc\", \"212\", \"1\", \"2\", \"3\", \"4\"),(\"def\", \"436\", \"2\", \"2\", \"1\", \"8\"),(\"abc\", \"510\", \"1\", \"2\", \"5\", \"8\")) val rdd = spark.sparkContext.parallelize(data) var df = spark.createDataFrame(rdd).toDF(columns:_*) df.show() for (cols <- df.columns.drop(2)) { df = df.withColumn(cols, foos_udf(col(\"C2\"),col(cols))) } df.show() for (cols <- df.columns.drop(2)) { df = df.withColumn(cols,col(cols).cast(FloatType)) } df.show()1,500 列数据错误:
Exception in thread \"main\" java.lang.StackOverflowError at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan.isStreaming(LogicalPlan.scala:37) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan$$anonfun$isStreaming$1.apply(LogicalPlan.scala:37) at org.apache.spark.sql.catalyst.plans.logical.LogicalPlan$$anonfun$isStreaming$1.apply(LogicalPlan.scala:37) at scala.collection.LinearSeqOptimized$class.exists(LinearSeqOptimized.scala:93) at scala.collection.immutable.List.exists(List.scala:84) ...想法:
- 也许
var可以替换,但数据大小接近RAM 的40%。 - 也许
for循环为dtype转换可能会导致性能下降,但我不知道如何,以及有哪些替代方案。通过在互联网上搜索,我看到有人建议基于foldLeft的方法,但显然仍然在内部转换为for循环。对此的任何投入将不胜感激。
- 也许
标签: scala apache-spark apache-spark-standalone