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)