avoid pivot when preprocessing data for vector assembler

Viewed 81

I have the pyspark dataframe below called id_predictions_df. I'm pivoting it to get it in the pivot_df2 form show below, so I can then run vectorassembler on it and get it in a form that I can perform k means clustering on. The pivot step is very slow. Does anyone know a faster alternative to pivot? Or is there a way I can avoid having to pivot the dataframe before feeding it to k means?

id_predictions_df.show()

+------+------+--------+
|id_x  |id_y  |adj_prob|
+------+------+--------+
|286724|286915|0.0659  |
|286724|295607|0.0681  |
|286724|322403|0.06544 |
+------+------+--------+



id_y_values=[x.id_y for x in id_predictions_df.select('id_y').distinct().collect()]

pivot_df2 = id_predictions_df.groupBy('id_x').pivot('id_y',id_y_values).agg({'adj_prob':'max'})



pivot_df2.show()


+------+------+--------------+
|id_x  |286915|295607| 322403|
+------+------+--------------+
|286724|0.0659|0.0681|0.06544|
+------+------+--------------+




from pyspark.ml.clustering import KMeans

from pyspark.ml.linalg import Vectors
from pyspark.ml.feature import VectorAssembler

feat_cols = [x for x in pivot_df2.columns if x!='id_x']

vec_assembler = VectorAssembler(inputCols = feat_cols, outputCol='features')

final_data = vec_assembler.transform(pivot_df2)

kmeans3 = KMeans(featuresCol='features',k=cluster_number)

model_k3 = kmeans3.fit(final_data)

cluster_label_df=model_k3.transform(final_data)
0 Answers
Related