如何使用PySpark将包含嵌套结构体数组的DataFrame展平至单层结构
展平嵌套PySpark DataFrame至单层结构
搞定这个嵌套DataFrame的扁平化其实很 straightforward,咱们一步步来拆解:
步骤1:导入必要函数
首先得把处理数组和列的核心函数导入进来:
from pyspark.sql.functions import explode, explode_outer, col
步骤2:先拆Books数组
你的Books是数组类型,我们先用explode把每个Book元素拆成单独的行——如果想保留那些Books为null的异常行,就换成explode_outer:
# 展开Books数组,生成新的Book结构体列,同时删掉原Books列 df_with_book = df.withColumn("Book", explode(col("Books"))).drop("Books")
步骤3:再拆Chapters数组
接下来处理Book结构体里的Chapters数组,同样用拆数组的方法:
# 这里用explode_outer是为了保留没有章节(Chapters为null)的书籍行 # 不需要这类行的话,直接换成explode就行 df_with_chapter = df_with_book.withColumn("Chapter", explode_outer(col("Book.Chapters"))).drop("Book.Chapters")
步骤4:展开结构体字段并重命名冲突列
现在所有数组都拆成行了,接下来把结构体里的字段全部展开,注意重命名重复的列名(比如作者的NAME和章节的NAME):
flattened_df = df_with_chapter.select( col("AUTHOR_ID"), col("NAME").alias("AUTHOR_NAME"), # 重命名避免和章节名冲突 col("Book.BOOK_ID"), col("Chapter.NAME").alias("CHAPTER_NAME"), col("Chapter.NUMBER_PAGES") )
最终的单层Schema
展平后的DataFrame会是这样的简洁结构:
root |-- AUTHOR_ID: integer (nullable = false) |-- AUTHOR_NAME: string (nullable = true) |-- BOOK_ID: integer (nullable = false) |-- CHAPTER_NAME: string (nullable = true) |-- NUMBER_PAGES: integer (nullable = true)
完整代码整合
把所有步骤串起来的完整示例:
from pyspark.sql.functions import explode_outer, col # 假设你的原始DataFrame名为df df_with_book = df.withColumn("Book", explode(col("Books"))).drop("Books") df_with_chapter = df_with_book.withColumn("Chapter", explode_outer(col("Book.Chapters"))).drop("Book.Chapters") flattened_df = df_with_chapter.select( "AUTHOR_ID", col("NAME").alias("AUTHOR_NAME"), "Book.BOOK_ID", col("Chapter.NAME").alias("CHAPTER_NAME"), "Chapter.NUMBER_PAGES" ) # 查看结果结构 flattened_df.printSchema()
内容的提问来源于stack exchange,提问作者Smaillns
相关产品推荐
相关产品推荐

