PySpark DataFrame写入S3指定目录结构问题求助
解决PySpark写入CSV到指定S3目录结构的问题
你的代码存在几个关键问题,以下是针对性的修正方案:
问题分析
- DataFrame不可变性忽略:
df.withColumn("datetime", current_timestamp())不会修改原DataFrame,必须将结果重新赋值给df才会生效。 - 时间格式不匹配:
current_timestamp()生成的是Timestamp类型,partitionBy会默认生成带时区的目录名(如datetime=2023-05-15T03:12:23.123+00:00),不符合你需要的20230515031223格式。 - 写入路径逻辑错误:Spark写入CSV时,指定的路径是目录而非文件名,直接把
input_file_actual.csv加到路径里会创建同名目录,而非目标文件。
修正步骤
1. 导入必要函数
先导入格式化时间的工具函数:
from pyspark.sql.functions import current_timestamp, date_format
2. 添加格式化后的datetime列
将当前时间转为yyyyMMddHHmmss格式的字符串,作为分区列:
# 生成符合要求的datetime字符串列 df = df.withColumn("datetime", date_format(current_timestamp(), "yyyyMMddHHmmss"))
3. 写入基础目录(自动生成分区结构)
指定基础路径为s3a://{bucket}/data/data,Spark会自动根据datetime列的值生成datetime=20230515031223目录:
bucket = "my-internal-bucket" base_path = f"s3a://{bucket}/data/data" # 写入CSV,开启表头,按datetime分区 df.write.option("header", "true").partitionBy("datetime").mode("overwrite").csv(base_path)
4. 重命名生成的part文件为指定文件名
Spark写入后会在分区目录下生成类似part-00000-xxx.csv的文件,需要将其重命名为input_file_actual.csv,可以用Hadoop文件系统API实现:
from pyspark.sql import SparkSession # 获取Hadoop文件系统实例 spark = SparkSession.builder.getOrCreate() fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration()) # 获取生成的datetime值,定位分区目录 datetime_value = df.select("datetime").first()[0] partition_dir = f"{base_path}/datetime={datetime_value}" target_file = f"{partition_dir}/input_file_actual.csv" # 筛选分区目录下的CSV文件(排除_SUCCESS标记文件) file_status = fs.listStatus(spark._jvm.org.apache.hadoop.fs.Path(partition_dir)) csv_files = [status.getPath().toString() for status in file_status if status.getPath().getName().endswith(".csv")] # 重命名第一个CSV文件为目标文件名,删除其他冗余文件 if csv_files: fs.rename(spark._jvm.org.apache.hadoop.fs.Path(csv_files[0]), spark._jvm.org.apache.hadoop.fs.Path(target_file)) for f in csv_files[1:]: fs.delete(spark._jvm.org.apache.hadoop.fs.Path(f), False)
可选:合并为单个文件再写入(小数据量适用)
如果数据量不大,可先合并为一个分区,避免生成多个part文件:
# 合并为1个分区写入临时目录 temp_path = f"{base_path}/temp_{datetime_value}" df.coalesce(1).write.option("header", "true").mode("overwrite").csv(temp_path) # 移动临时文件到目标路径并重命名 temp_csv_files = [status.getPath().toString() for status in fs.listStatus(spark._jvm.org.apache.hadoop.fs.Path(temp_path)) if status.getPath().getName().endswith(".csv")] if temp_csv_files: fs.rename(spark._jvm.org.apache.hadoop.fs.Path(temp_csv_files[0]), spark._jvm.org.apache.hadoop.fs.Path(target_file)) # 删除临时目录 fs.delete(spark._jvm.org.apache.hadoop.fs.Path(temp_path), True)
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import current_timestamp, date_format # 初始化SparkSession(未初始化时执行) spark = SparkSession.builder.appName("WriteCSVToS3").getOrCreate() # 读取原DataFrame(你的原有代码) bucket = "my-internal-bucket" directory_remote = "data" toName_csv = "your_input_file.csv" schema = your_schema_definition # 替换为你的实际Schema df = spark.read.format('csv').options(header='true').load(f"s3a://{bucket}/{directory_remote}/{toName_csv}", schema=schema) # 添加格式化后的datetime列 df = df.withColumn("datetime", date_format(current_timestamp(), "yyyyMMddHHmmss")) # 定义路径变量 base_path = f"s3a://{bucket}/data/data" datetime_value = df.select("datetime").first()[0] target_file_path = f"{base_path}/datetime={datetime_value}/input_file_actual.csv" # 合并分区写入临时目录 temp_path = f"{base_path}/temp_{datetime_value}" df.coalesce(1).write.option("header", "true").mode("overwrite").csv(temp_path) # 处理文件重命名与清理 fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration()) temp_csv_files = [status.getPath().toString() for status in fs.listStatus(spark._jvm.org.apache.hadoop.fs.Path(temp_path)) if status.getPath().getName().endswith(".csv")] if temp_csv_files: fs.rename(spark._jvm.org.apache.hadoop.fs.Path(temp_csv_files[0]), spark._jvm.org.apache.hadoop.fs.Path(target_file_path)) fs.delete(spark._jvm.org.apache.hadoop.fs.Path(temp_path), True) print(f"文件已成功写入:{target_file_path}")
注意:coalesce(1)会将所有数据集中到单个Executor,不适合超大数据量场景。如果数据量较大,建议保留多part文件,或使用合理的分区策略拆分数据,避免单文件过大。
内容的提问来源于stack exchange,提问作者marie20
相关产品推荐
相关产品推荐

