Spark技术问询:能否将分区列包含在输出文件内容中?
如何让Spark分区写入CSV时包含分区列
当然可以实现这个需求!Spark默认会把分区列只放在路径目录中,不会写入到CSV文件里,但我们有几种简单的办法来让foo列也出现在CSV内容中,下面给你详细说明:
方法一:复制分区列后再写入
因为Spark在使用partitionBy时会自动将分区列从数据文件中排除,所以我们可以先给DataFrame新增一个和foo值完全相同的列(比如命名为foo_val),这样写入时这个新增列会被保留在CSV里,同时依然按foo字段分区:
PySpark示例
from pyspark.sql.functions import col # 新增一个与foo值一致的列 df_with_foo = df.withColumn("foo_val", col("foo")) # 按foo分区写入CSV,此时文件内会包含bar和foo_val列 df_with_foo.write.mode("overwrite").format("csv").partitionBy("foo").save("/tmp/output")
Scala示例
import org.apache.spark.sql.functions.col val dfWithFoo = df.withColumn("foo_val", col("foo")) dfWithFoo.write.mode("overwrite").format("csv").partitionBy("foo").save("/tmp/output")
方法二:手动分区写入(不依赖partitionBy)
如果你不想新增列,希望CSV里直接保留原始的foo列,可以放弃使用partitionBy,手动遍历所有唯一的foo值,过滤对应数据后写入指定的分区目录:
PySpark示例
from pyspark.sql.functions import col # 获取所有唯一的foo值 foo_list = [row.foo for row in df.select("foo").distinct().collect()] for foo_val in foo_list: # 筛选出当前foo值对应的数据集 target_df = df.filter(col("foo") == foo_val) # 写入到对应分区目录,生成的CSV会包含foo和bar列 target_df.write.mode("overwrite").format("csv").save(f"/tmp/output/foo={foo_val}")
两种方法各有优劣:方法一更简洁高效,适合大多数常规场景;方法二更灵活,适合需要对每个分区数据做额外自定义处理的情况。
内容的提问来源于stack exchange,提问作者erwaman
相关产品推荐
相关产品推荐

