Scala读取多行JSON,从含嵌套结构的DataFrame提取指定30列
从嵌套结构Spark DataFrame提取指定列的方法
嘿,这个需求我熟!针对这种嵌套了数组+结构体的Spark DataFrame,提取指定字段其实分两种常见场景,我给你一步步讲清楚:
场景1:把数组展开,每个结构体元素拆成单独行
如果你的需求是把Applications数组里的每个结构体元素都拆成独立的行,同时提取需要的30个字段(包括顶级列和嵌套字段),可以这么做:
步骤:
- 先用
explode函数把数组列展开,将每个数组元素转化为单独的行; - 用
.操作符直接提取嵌套结构体里的字段,再搭配select选择所有需要的列。
代码示例:
import org.apache.spark.sql.functions.explode // 第一步:展开Applications数组,得到每个结构体的单独行 val explodedDF = peopleDF.withColumn("app", explode(col("Applications"))) // 第二步:选择需要的字段——这里假设你有顶级列(比如user_id),再加上嵌套字段 val selectedDF = explodedDF.select( col("user_id"), // 替换成你实际需要的顶级列 col("app.b_als_o_isehp").alias("als_o_isehp"), // 重命名字段让结果更清晰 col("app.b_als_p_isehp").alias("als_p_isehp"), col("app.l_als_o_eventid").alias("als_o_eventid"), // ... 这里继续添加剩下的27个需要的字段,格式和上面一致 )
如果需要提取的字段太多(30个),一个个写太麻烦,可以把字段名整理成列表,动态生成选择表达式:
// 定义需要的字段列表,嵌套字段用"app.字段名"的格式 val requiredColumns = List( "user_id", "app.b_als_o_isehp", "app.b_als_p_isehp", "app.l_als_o_eventid", // ... 补充剩下的字段 ) // 动态选择所有列 val selectedDF = explodedDF.select(requiredColumns.map(col): _*)
场景2:保留数组结构,只提取结构体里的指定字段
如果你不想拆分数组,只想让数组里的每个结构体只保留你需要的字段,那可以用transform函数来处理数组元素:
代码示例:
import org.apache.spark.sql.functions.transform // 定义需要从结构体里提取的字段列表 val nestedRequiredFields = List( "b_als_o_isehp", "b_als_p_isehp", "l_als_o_eventid", // ... 补充其他嵌套字段 ) // 用transform遍历数组,每个结构体只保留需要的字段 val transformedDF = peopleDF.withColumn( "filtered_applications", transform( col("Applications"), app => struct( // 遍历字段列表,逐个提取并重命名(可选) nestedRequiredFields.map(field => app.getField(field).alias(field)): _* ) ) ) // 最后选择需要的顶级列 + 处理后的数组列 val finalDF = transformedDF.select( col("some_top_level_column"), // 替换成实际顶级列 col("filtered_applications") // ... 补充其他需要的顶级列 )
注意事项:
- 如果字段名包含特殊字符(比如空格、特殊符号),记得用反引号包裹,比如
col("app.b als o isehp"); - Spark默认字段名不区分大小写,但开启严格模式时要注意大小写匹配;
- 原Schema里的字段都是nullable的,提取后依然保持nullable,如果需要处理空值,可以用
coalesce等函数填充默认值。
内容的提问来源于stack exchange,提问作者Big data Hadoop dev.
相关产品推荐
相关产品推荐

