如何让Spark按tenantId分区生成无列名前缀的目录结构?
解决Spark分区目录不带列名的方案
Spark默认的partitionBy会生成列名=值格式的分区目录,要实现直接用分区值作为目录名(如C:\1\sample1.csv),可以用以下两种方法:
方法一:按分区值分组后逐个写入
这种方法适合中小规模数据集,直接遍历每个唯一的tenantId,过滤数据后写入指定目录,同时将输出文件重命名为sample1.csv。
Scala 代码示例
import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.hadoop.conf.Configuration // 读取CSV数据 val df = spark.read.csv("path/to/your/input.csv").toDF("tenantId", "siteId", "TrackingID") // 获取所有唯一的tenantId val tenantIds = df.select("tenantId").distinct().collect().map(_.getString(0)) // 遍历每个tenantId,写入并整理文件 tenantIds.foreach { id => val tempDir = s"C:\\temp_$id" // 过滤当前tenantId的数据,合并为一个文件写入临时目录 df.filter(s"tenantId = '$id'") .coalesce(1) .write .mode("overwrite") .csv(tempDir) // 初始化Hadoop文件系统 val fs = FileSystem.get(new Configuration()) val tempPath = new Path(tempDir) // 找到临时目录下的part文件 val partFile = fs.listStatus(tempPath) .filter(status => status.isFile && status.getPath.getName.startsWith("part-")) .head.getPath // 目标文件路径 val destDir = new Path(s"C:\\$id") val destFile = new Path(destDir, "sample1.csv") // 创建目标目录(如果不存在) if (!fs.exists(destDir)) fs.mkdirs(destDir) // 删除已存在的目标文件 if (fs.exists(destFile)) fs.delete(destFile, false) // 重命名part文件为sample1.csv fs.rename(partFile, destFile) // 删除临时目录 fs.delete(tempPath, true) }
Python 代码示例
from pyspark.sql import SparkSession import os import shutil spark = SparkSession.builder.appName("CustomPartition").getOrCreate() # 读取带表头的CSV数据 df = spark.read.csv("path/to/your/input.csv", header=True) # 获取所有唯一的tenantId tenant_ids = [row.tenantId for row in df.select("tenantId").distinct().collect()] # 遍历每个tenantId处理数据 for tenant_id in tenant_ids: temp_dir = f"C:\\temp_{tenant_id}" # 过滤数据并合并为单个文件写入临时目录 df.filter(df.tenantId == tenant_id) \ .coalesce(1) \ .write \ .mode("overwrite") \ .csv(temp_dir) # 查找临时目录下的part文件(排除crc校验文件) part_file = None for file_name in os.listdir(temp_dir): if file_name.startswith("part-") and not file_name.endswith(".crc"): part_file = os.path.join(temp_dir, file_name) break # 准备目标目录和文件 dest_dir = f"C:\\{tenant_id}" dest_file = os.path.join(dest_dir, "sample1.csv") # 创建目标目录 os.makedirs(dest_dir, exist_ok=True) # 删除已存在的目标文件 if os.path.exists(dest_file): os.remove(dest_file) # 移动并重命名文件 shutil.move(part_file, dest_file) # 删除临时目录 shutil.rmtree(temp_dir)
方法二:先按默认分区写入,再重命名目录
这种方法适合大规模数据集,先利用Spark的partitionBy高效写入,再通过文件系统API批量重命名分区目录,并整理文件名。
Scala 代码示例
import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.hadoop.conf.Configuration // 读取CSV数据 val df = spark.read.csv("path/to/your/input.csv").toDF("tenantId", "siteId", "TrackingID") // 临时输出目录 val tempOutputDir = "C:\\temp_partition" // 用默认partitionBy写入数据 df.write.partitionBy("tenantId").mode("overwrite").csv(tempOutputDir) // 初始化Hadoop文件系统 val fs = FileSystem.get(new Configuration()) val basePath = new Path(tempOutputDir) // 遍历所有分区目录 fs.listStatus(basePath).filter(_.isDirectory).foreach { status => val partitionDir = status.getPath val dirName = partitionDir.getName // 格式为tenantId=1 val tenantId = dirName.split("=")(1) // 目标目录 val targetDir = new Path(s"C:\\$tenantId") // 删除已存在的目标目录 if (fs.exists(targetDir)) fs.delete(targetDir, true) // 重命名分区目录 fs.rename(partitionDir, targetDir) // 将目录下的part文件重命名为sample1.csv val partFile = fs.listStatus(targetDir) .filter(s => s.isFile && s.getPath.getName.startsWith("part-")) .head.getPath val targetFile = new Path(targetDir, "sample1.csv") if (fs.exists(targetFile)) fs.delete(targetFile, false) fs.rename(partFile, targetFile) } // 删除临时目录 fs.delete(basePath, true)
Python 代码示例
from pyspark.sql import SparkSession import os import shutil spark = SparkSession.builder.appName("RenamePartition").getOrCreate() # 读取带表头的CSV数据 df = spark.read.csv("path/to/your/input.csv", header=True) # 临时输出目录 temp_output_dir = "C:\\temp_partition" # 用默认partitionBy写入数据 df.write.partitionBy("tenantId").mode("overwrite").csv(temp_output_dir) # 遍历临时目录下的分区目录 for dir_name in os.listdir(temp_output_dir): if dir_name.startswith("tenantId="): tenant_id = dir_name.split("=")[1] src_dir = os.path.join(temp_output_dir, dir_name) dest_dir = f"C:\\{tenant_id}" # 删除已存在的目标目录 if os.path.exists(dest_dir): shutil.rmtree(dest_dir) # 重命名分区目录 shutil.move(src_dir, dest_dir) # 将目录下的part文件重命名为sample1.csv part_file = None for file_name in os.listdir(dest_dir): if file_name.startswith("part-") and not file_name.endswith(".crc"): part_file = os.path.join(dest_dir, file_name) break dest_file = os.path.join(dest_dir, "sample1.csv") if os.path.exists(dest_file): os.remove(dest_file) os.rename(part_file, dest_file) # 删除临时目录 shutil.rmtree(temp_output_dir)
注意事项
- 使用
coalesce(1)会将分区数据合并为一个文件,适合需要单个sample1.csv的场景,但如果数据量极大,可能会导致单个文件过大,此时可以考虑保留多个part文件,或者根据实际情况调整合并的分区数。 - 两种方法都需要注意目录权限问题,确保Spark程序有读写目标目录的权限。
- 如果是在分布式文件系统(如HDFS)上操作,方法二的效率会更高,因为Spark的
partitionBy能利用分布式写入的优势。
内容的提问来源于stack exchange,提问作者venkatesh bandaru
相关产品推荐
相关产品推荐

