PySpark中写入Parquet文件时如何保留原文件名?
Spark写入Parquet时能否保留原自定义文件名?
答案是默认情况下不行。
Spark的df.write.parquet()方法有固定的文件命名逻辑,写入时会自动生成类似part-00000.snappy.parquet这类标准格式的文件,完全不会读取或保留你原目录中fr_default_players_results.parquet、us_default_players_results.parquet这类自定义文件名。
如果一定要保留原文件名,只能通过手动处理单个分区文件的方式实现,但这种方式会牺牲Spark的分布式处理优势,仅适合小数据量场景:
- 先遍历原分区目录,获取每个自定义Parquet文件的完整路径和对应的目标文件名
- 针对每个文件单独用Spark读取、处理
- 将处理后的DataFrame用
coalesce(1)合并为单个分区(避免生成多个part文件),写入到临时目录 - 最后通过文件系统API(比如Hadoop的
FileSystem或本地文件操作)把临时目录里的part文件重命名为原自定义文件名,移动到目标分区路径下
举个简化的Scala代码示例:
import org.apache.hadoop.fs.{FileSystem, Path} // 获取文件系统实例 val fs = FileSystem.get(spark.sparkContext.hadoopConfiguration) // 原数据根目录 val rootPath = new Path("/path/to/Users_data") // 遍历所有自定义Parquet文件 val parquetFiles = fs.globStatus(new Path(rootPath, "**/*.parquet")).map(_.getPath) parquetFiles.foreach { filePath => // 读取单个文件 val df = spark.read.parquet(filePath.toString) // 执行数据处理逻辑 val processedDf = df.filter("score > 80") // 临时输出目录 val tempDir = new Path(filePath.getParent, "temp_" + filePath.getName) // 合并为单个分区写入临时目录 processedDf.coalesce(1).write.mode("overwrite").parquet(tempDir.toString) // 找到临时目录里的part文件 val partFile = fs.globStatus(new Path(tempDir, "part*.parquet")).head.getPath // 重命名为原文件名并覆盖原位置 fs.rename(partFile, filePath) // 删除临时目录 fs.delete(tempDir, true) }
需要注意的是,这种方法不适合大数据量:每个文件单独处理无法利用Spark的分布式并行能力,coalesce(1)还会把数据集中到单个节点,可能引发内存溢出问题。
内容的提问来源于stack exchange,提问作者Yefet
相关产品推荐
相关产品推荐

