PySpark读取JSON转为Dataframe并提取嵌套数组为目标列的实现方法
PySpark 实现方案
共有两种常用的实现方式,优先推荐使用高阶函数的方案,性能更优。
方案1:使用transform + concat_ws 高阶函数(推荐)
无需拆分行,直接在原地对数组字段做处理,适合所有数据规模,Spark 3.1+支持lambda语法调用,更低版本可以用expr写法实现。
完整代码
from pyspark.sql import SparkSession from pyspark.sql.functions import transform, concat_ws, col, expr # 初始化SparkSession spark = SparkSession.builder.appName("ParseJsonSample").getOrCreate() # 读取JSON文件,自动识别JSON结构 df = spark.read.json("你的JSON文件路径.json") # Spark 3.1+ 写法 result_df = df.withColumn( "sample", concat_ws( ",", transform(col("sample"), lambda x: x.getField("value")) ) ) # Spark 3.1以下版本写法,替换上面的withColumn逻辑即可 # result_df = df.withColumn( # "sample", # expr("concat_ws(',', transform(sample, item -> item.value))") # ) # 输出结果 result_df.show(truncate=False)
运行后输出结果和你预期的格式完全一致。
方案2:使用explode + group by 聚合
兼容性更好,支持所有Spark版本,但会先把数组拆分为多行再聚合,大数据量下性能略差。
完整代码
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, collect_list, concat_ws, col spark = SparkSession.builder.appName("ParseJsonSample").getOrCreate() df = spark.read.json("你的JSON文件路径.json") result_df = df.withColumn("sample_item", explode(col("sample"))) \ .groupBy("scientific_names", "gender", "examples") \ .agg(concat_ws(",", collect_list(col("sample_item.value"))).alias("sample")) \ .select("sample", "scientific_names", "gender", "examples") result_df.show(truncate=False)
内容的提问来源于stack exchange,提问作者ben
相关产品推荐
相关产品推荐

