使用PySpark拆分单行到多行:explode空值问题与示例请求
PySpark 实现多描述列转多行并过滤空值
原始数据集
| ID | Name | Desc1 | Desc2 | Desc3 | Desc4 |
|---|---|---|---|---|---|
| 100 | ABC | A | B | C | D |
| 200 | XYZ | X | Y |
期望输出
| ID | Name | Desc |
|---|---|---|
| 100 | ABC | A |
| 100 | ABC | B |
| 100 | ABC | C |
| 100 | ABC | D |
| 200 | XYZ | X |
| 200 | XYZ | Y |
问题分析
你尝试用explode实现转置,但代码存在多处语法错误(列名访问错误、变量名错误、括号缺失),且未处理原始数据中的空字符串,导致结果出现空值行。
正确实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import array, explode, col, filter # 初始化SparkSession spark = SparkSession.builder.appName("UnpivotDescriptions").getOrCreate() # 构建带Schema的DataFrame dept_data = [(100,"ABC","A","B","C","D"),(200,"XYZ","X","Y","","")] schema_cols = ["ID","NAME","Desc1","Desc2","Desc3","Desc4"] df = spark.createDataFrame(dept_data, schema=schema_cols) # 合并Desc列成数组,并过滤空字符串和null desc_columns = [col(col_name) for col_name in schema_cols if col_name.startswith("Desc")] filtered_desc_array = filter(array(*desc_columns), lambda x: x != "" and x.isNotNull()) # 展开数组得到最终结果 result_df = df.select("ID", "NAME", explode(filtered_desc_array).alias("Desc")) # 输出结果 result_df.show()
关键说明
- Schema规范:直接用
createDataFrame指定列名,避免RDD转换时的列名混乱,确保后续列访问准确 - 空值过滤:使用Spark的
filter函数(3.0+版本支持)移除数组中的空字符串和null,避免explode后产生无效空行 - 动态列选择:通过
startswith("Desc")动态匹配所有描述列,代码更灵活,后续新增Desc列无需修改逻辑
运行结果
+---+----+----+ | ID|NAME|Desc| +---+----+----+ |100| ABC| A| |100| ABC| B| |100| ABC| C| |100| ABC| D| |200| XYZ| X| |200| XYZ| Y| +---+----+----+
内容的提问来源于stack exchange,提问作者raw
相关产品推荐
相关产品推荐

