在Databricks Python Notebook中拆分MongoDB生成的Spark DataFrame数组为多行
嘿,这个需求在Spark里其实挺常见的,用explode()系列函数就能轻松搞定!我给你一步步拆解操作方法,还有不同场景的应对方案:
基础操作:拆分数组为多行
假设你已经从MongoDB读取并创建了Spark DataFrame(比如叫mongo_df),其中有一个数组类型的字段(比如skills)需要拆分。
步骤1:导入所需函数
首先需要导入Spark的explode函数:
from pyspark.sql.functions import explode
步骤2:执行拆分操作
使用select方法,保留原DataFrame的所有列,同时将数组字段拆分为多行:
# 拆分skills数组,新增一个拆分后的元素列skill exploded_df = mongo_df.select("*", explode("skills").alias("skill"))
举个实际例子:
原DataFrame(
mongo_df)的数据:
user_id username skills 101 zhangsan ["Python", "Spark"] 102 lisi ["Java"]
执行拆分后,exploded_df会变成:
| user_id | username | skills | skill |
|---|---|---|---|
| 101 | zhangsan | ["Python", "Spark"] | Python |
| 101 | zhangsan | ["Python", "Spark"] | Spark |
| 102 | lisi | ["Java"] | Java |
如果你不想保留原来的数组列,可以用drop去掉它:
exploded_df = mongo_df.drop("skills").select("*", explode("skills").alias("skill"))
进阶场景处理
1. 保留空数组/Null的行
默认的explode()会过滤掉数组为空或者为Null的行,如果想保留这些行,可以用explode_outer():
from pyspark.sql.functions import explode_outer exploded_df = mongo_df.select("*", explode_outer("skills").alias("skill"))
这样数组为空的行依然会保留,拆分后的skill列值为Null。
2. 同时获取数组元素的索引
如果需要知道每个拆分元素在原数组中的位置,可以用posexplode()(或者posexplode_outer()保留空值):
from pyspark.sql.functions import posexplode # 同时获取索引和元素,分别命名为skill_index和skill exploded_df = mongo_df.select("*", posexplode("skills").alias("skill_index", "skill"))
拆分后的结果会多一列索引值:
| user_id | username | skills | skill_index | skill |
|---|---|---|---|---|
| 101 | zhangsan | ["Python", "Spark"] | 0 | Python |
| 101 | zhangsan | ["Python", "Spark"] | 1 | Spark |
| 102 | lisi | ["Java"] | 0 | Java |
内容的提问来源于stack exchange,提问作者Raju
相关产品推荐
相关产品推荐

