如何在PySpark中按指定多条件生成新列A5并拆分对应行?
实现Spark DataFrame条件行拆分与新列生成
核心思路
先通过条件判断生成包含所有符合要求的a5值的数组列,再利用explode函数将数组拆分为多行,每个数组元素对应一行数据。
步骤与代码示例(PySpark)
- 导入依赖函数
from pyspark.sql import SparkSession from pyspark.sql import functions as F
- 构造示例DataFrame
# 初始化SparkSession spark = SparkSession.builder.appName("ConditionRowSplit").getOrCreate() # 构造原始数据 data = [("A", 12, 9, 1), ("B", 14, 13, 1), ("C", 7, 3, 0)] df = spark.createDataFrame(data, schema=["a1", "a2", "a3", "a4"])
- 生成包含符合条件的a5值的数组列
通过array收集每个条件对应的结果,不满足条件的会返回null,再用filter剔除数组中的null值:
df_with_array = df.withColumn( "a5_array", F.filter( F.array( F.when(F.col("a1") == "A", "Car"), # 条件1:a1=A时对应Car F.when(F.col("a2") > 0, "Bus"), # 条件2:a2>0时对应Bus F.when((F.col("a3") > 0) & (F.col("a4") == 1), "Bike") # 条件3:a3>0且a4=1时对应Bike ), lambda x: x.isNotNull() ) )
- 拆分数组列得到最终结果
用explode函数将数组列拆分为多行,每个元素对应一行:
result_df = df_with_array.select("a1", "a2", "a3", "a4", F.explode("a5_array").alias("a5")) # 查看结果 result_df.show()
代码示例(Scala)
如果使用Scala版本,代码逻辑一致:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ object ConditionRowSplit { def main(args: Array[String]): Unit = { val spark = SparkSession.builder.appName("ConditionRowSplit").getOrCreate() val data = Seq(("A", 12, 9, 1), ("B", 14, 13, 1), ("C", 7, 3, 0)) val df = data.toDF("a1", "a2", "a3", "a4") val dfWithArray = df.withColumn( "a5_array", filter( array( when(col("a1") === "A", "Car"), when(col("a2") > 0, "Bus"), when((col("a3") > 0) && (col("a4") === 1), "Bike") ), x => x.isNotNull ) ) val resultDF = dfWithArray.select($"a1", $"a2", $"a3", $"a4", explode($"a5_array").alias("a5")) resultDF.show() } }
结果验证
执行代码后,输出的DataFrame与期望完全一致:
| a1 | a2 | a3 | a4 | a5 |
|---|---|---|---|---|
| A | 12 | 9 | 1 | Car |
| A | 12 | 9 | 1 | Bus |
| A | 12 | 9 | 1 | Bike |
| B | 14 | 13 | 1 | Bus |
| B | 14 | 13 | 1 | Bike |
| C | 7 | 3 | 0 | Bus |
内容的提问来源于stack exchange,提问作者Sadhvi Lakshmi
相关产品推荐
相关产品推荐

