【发布时间】:2018-05-08 04:09:03
【问题描述】:
我有一个相当复杂的过程来创建 pyspark 数据框,将其转换为 pandas 数据框,并将结果输出到平面文件。不知道是在什么时候引入了错误,所以我将描述整个过程。
一开始我有一个 pyspark 数据框,其中包含 id 集的成对相似性。它看起来像这样:
+------+-------+-------------------+
| ID_A| ID_B| EuclideanDistance|
+------+-------+-------------------+
| 1| 1| 0.0|
| 1| 2|0.13103884200454394|
| 1| 3| 0.2176246463836219|
| 1| 4| 0.280568636550471|
...
我喜欢按 ID_A 对其进行分组,按 EuclideanDistance 对每个组进行排序,并且只获取每个组的前 N 对。所以首先我这样做:
from pyspark.sql.window import Window
from pyspark.sql.functions import rank, col, row_number
window = Window.partitionBy(df['ID_A']).orderBy(df_sim['EuclideanDistance'])
result = (df.withColumn('row_num', row_number().over(window)))
我确保 ID_A = 1 仍在“结果”数据框中。然后我这样做是为了将每个组限制为 20 行:
result1 = result.where(result.row_num<20)
result1.toPandas().to_csv("mytest.csv")
并且 ID_A = 1 不在生成的 .csv 文件中(尽管它仍然存在于 result1 中)。在这个转换链中的某个地方是否存在可能导致数据丢失的问题?
【问题讨论】: