PySpark如何拆分含结构体的数组字段并导出为CSV
PySpark 动态拆分嵌套数组为多列实现方案
核心实现步骤
- 预处理嵌套数组:将每个结构体的
lang1、lang2字段用/拼接为单个字符串,生成新的字符串数组字段 - 动态计算数组最大长度:无需硬编码列数,自动适配不同长度的数组
- 按规则生成多列:根据最大长度生成对应数量的
langage_N列,不足长度的位置自动填充null
完整可运行代码
from pyspark.sql import functions as F # 第一步:将数组内结构体的两个字段拼接为字符串,生成新数组字段 # 低版本Spark可替换为expr写法:F.expr("transform(langages, x -> concat_ws('/', x.lang1, x.lang2))") df_processed = df.withColumn( "lang_concat", F.transform("langages", lambda x: F.concat_ws("/", x.lang1, x.lang2)) ) # 第二步:计算langages数组的最大长度,动态确定生成列的数量 max_array_len = df_processed.agg(F.max(F.size("lang_concat"))).head()[0] # 第三步:生成指定格式的列,选择需要导出的字段 output_cols = [F.col("firstname"), F.col("lastname")] + [ F.element_at("lang_concat", idx).alias(f"langage_{idx}") for idx in range(1, max_array_len + 1) ] df_result = df_processed.select(output_cols) # 导出为CSV df_result.write.csv("./lang_output", header=True, encoding="utf-8")
输出效果说明
基于你提供的测试数据,处理后的结果如下:
+---------+--------+------------+-----------+-------------+ |firstname|lastname| langage_1 | langage_2 | langage_3 | +---------+--------+------------+-----------+-------------+ | john| smith| Java/Python| C/R| Perl/Scala | | robert| plant| C/Java|Python/Perl| null| +---------+--------+------------+-----------+-------------+
数组长度不足的位置会自动填充null,导出CSV时默认展示为空值,符合CSV格式规范。
内容的提问来源于stack exchange,提问作者Fabrice
相关产品推荐
相关产品推荐

