【发布时间】:2017-01-22 12:09:30
【问题描述】:
我正在准备一个带有 id 和我的特征向量的 DataFrame,以便稍后用于进行预测。我在我的数据框上做了一个 groupBy,在我的 groupBy 中,我将几列作为列表合并到一个新列中:
def mergeFunction(...) // with 14 input variables
val myudffunction( mergeFunction ) // Spark doesn't support this
df.groupBy("id").agg(
collect_list(df(...)) as ...
... // too many of these (something like 14 of them)
).withColumn("features_labels",
myudffunction(
col(...)
, col(...) )
.select("id", "feature_labels")
这就是我创建特征向量及其标签的方式。到目前为止它一直对我有用,但这是我的特征向量第一次使用这种方法变得大于数字 10,这是 Spark 中的 udf 函数最多接受的。
我不确定我还能如何解决这个问题?是 udf 输入的大小 Spark 会变大,我是不是理解错了,或者 有更好的方法吗?
【问题讨论】:
标签: scala apache-spark dataframe apache-spark-sql apache-spark-mllib