含Array与Struct的Parquet表转CSV的Spark处理方案咨询
Parquet嵌套结构(Array/Struct)转CSV的正确姿势
关键操作说明
要把带嵌套的Parquet转成扁平化的CSV,得先搞定两个核心点:
- Struct字段提取:直接用
.点符号钻取子字段即可,比如merchant.id就是取merchant结构体里的id字段; - Array字段展开:数组类型必须先用
explode函数拆成多行,不然CSV里会存数组的字符串形式,完全没法正常使用。
完整代码示例(Spark Scala)
先模拟你的表结构(如果实际结构有差异,对应调整字段名即可):
root |-- region: string |-- id: string |-- article: string |-- name: string |-- url: string |-- variations: array | |-- element: struct | | |-- msrp: double | | |-- price: double |-- merchant: struct | |-- id: string | |-- name: string
代码实现:
import org.apache.spark.sql.functions.explode // 读取Parquet表 val df = spark.read.parquet("/your/parquet/file/path") // 扁平化处理嵌套结构 val flattenedDf = df // 把variations数组拆成单行,每个数组元素对应一行数据 .withColumn("variation", explode(col("variations"))) // 选择目标字段,给嵌套字段起别名避免列名混乱 .select( "region", "id", "article", "name", "url", col("variation.msrp").alias("variation_msrp"), col("variation.price").alias("variation_price"), col("merchant.id").alias("merchant_id"), col("merchant.name").alias("merchant_name") ) // 保存为带表头的CSV flattenedDf.write .option("header", "true") .mode("overwrite") .csv("/your/csv/output/path")
补充注意事项
- 如果你的
variations数组可能为空,用explode_outer替代explode,这样不会丢掉原表中数组为空的行; - 必须给嵌套字段加别名(
alias),不然CSV列名会是variation.msrp这种格式,后续处理极易出错; - 用PySpark的话,逻辑完全一致,代码写法稍作调整:
from pyspark.sql.functions import explode, col df = spark.read.parquet("/your/parquet/file/path") flattened_df = df \ .withColumn("variation", explode(col("variations"))) \ .select( "region", "id", "article", "name", "url", col("variation.msrp").alias("variation_msrp"), col("variation.price").alias("variation_price"), col("merchant.id").alias("merchant_id"), col("merchant.name").alias("merchant_name") ) flattened_df.write \ .option("header", True) \ .mode("overwrite") \ .csv("/your/csv/output/path")
内容的提问来源于stack exchange,提问作者Kristina Isaenko
相关产品推荐
相关产品推荐

