【发布时间】:2020-02-06 20:32:42
【问题描述】:
我目前正在清理数据集,并且一直在尝试使用 pyspark。数据从 csv 读入数据帧,我需要的值在它们各自的行中,但对于某些行,值是混合的。我需要旋转这些行的值,以便这些值位于正确的列中。例如,假设我有以下数据集:
+-------+-------+-------+
| A | B | C |
+-------+-------+-------+
| 2 | 3 | 1 |
+-------+-------+-------+
但第一行的值应该是
+-------+-------+-------+
| A | B | C |
+-------+-------+-------+
| 1 | 2 | 3 |
+-------+-------+-------+
我当前的解决方案是添加一个临时列,并为每一列重新分配值,并在删除旧列的同时重命名临时列:
// Add temporary column C
+-------+-------+-------+-------+
| A | B | C | tmp_C |
+-------+-------+-------+-------+
| 2 | 3 | 1 | 1 |
+-------+-------+-------+-------+
// Shift values
+-------+-------+-------+-------+
| A | B | C | tmp_C |
+-------+-------+-------+-------+
| 2 | 2 | 3 | 1 |
+-------+-------+-------+-------+
// Drop old column
+-------+-------+-------+
| B | C | tmp_C |
+-------+-------+-------+
| 2 | 3 | 1 |
+-------+-------+-------+
// Rename new column
+-------+-------+-------+
| B | C | A |
+-------+-------+-------+
| 2 | 3 | 1 |
+-------+-------+-------+
我在 pyspark 中实现的方式如下:
from pyspark.sql import SparkSession
from pyspark.sql.function import when, col
def clean_data(spark_session, file_path):
df = (
spark_session
.read
.csv(file_path, header='true')
)
df = (
df
.withColumn(
"tmp_C",
when(
col("C") == 1,
col("C")
).otherwise("A")
)
.withColumn(
"C",
when(
col("C") == 1,
col("B")
).otherwise("C")
)
.withColumn(
"B",
when(
col("C") == 1,
col("A")
).otherwise("B")
)
)
df = df.drop("A")
df = df.withColumnRenamed("tmp_C", "A")
return df
对我来说,这看起来不太好,我不确定这是解决这个问题的最佳方法。我对 Spark 很陌生,想知道解决这种情况的最佳方法,尽管这确实有效。另外,我还想知道这是否是 Spark 的一个很好的用例(请注意,我使用的数据集很大,而且还有比这更多的字段。上面的例子大大简化了)。
【问题讨论】:
标签: pyspark