PySpark分块保存CSV后zip方法失效,OFS路径适配咨询
问题描述
我尝试创建一个类,从数据表中获取数据,按行分割为每个包含100万行的块,用PySpark将这些分区数据保存为独立CSV文件。目前已成功将文件保存到指定output_path,但自定义的zip_csv_files方法无法正常工作。我们使用ofs://oc/ectt这类OFS路径,请问是否需要在zip_csv_files方法中更换路径处理库?
执行代码
script = spark.sql('''select status_date, sum(ros) as ros_all from table_name group by 1 order by 1''') output_path = 'xxx' file_name_prefix = 'marikina' chunk_size = 1000000 saver = DataFrameChunkSaver(script, output_path, file_name_prefix, chunk_size) saver.save_chunks()
执行结果
save_chunks() 执行结果
- Chunk 1保存至
xxx/marikina0.folder - Chunk 2保存至
xxx/marikina1.folder
zip_csv_files() 执行结果
Zipping CSV files to xxx/chunks.zip Checking folder: xxx/marikina0.folder Folder xxx/marikina0.folder does not exist or is not a directory. Stopping.
相关类代码
import os import zipfile class DataFrameChunkSaver: def __init__(self, df, output_path, file_name_prefix='karl', chunk_size=1000000): self.df = df self.output_path = output_path self.chunk_size = chunk_size self.file_name_prefix = file_name_prefix self.output_file_template = os.path.join(self.output_path, f"{self.file_name_prefix}{{}}.folder") self.zip_file_path = os.path.join(self.output_path, "chunks.zip") # Ensure the directory exists os.makedirs(self.output_path, exist_ok=True) def count_chunks(self): total_rows = self.df.count() return -(-total_rows // self.chunk_size) def repartition_df(self): num_chunks = self.count_chunks() return self.df.repartition(num_chunks) def save_chunks(self): repartitioned_df = self.repartition_df() num_chunks = self.count_chunks() for i in range(num_chunks): chunk_file_path = self.output_file_template.format(i) partition_df = repartitioned_df.filter(f"spark_partition_id() == {i}") partition_df.coalesce(1).write.csv(chunk_file_path, header=True) print(f"Chunk {i + 1} saved to {chunk_file_path}")
def zip_csv_files(self): print(f"Zipping CSV files to {self.zip_file_path}") with zipfile.ZipFile(self.zip_file_path, 'w', zipfile.ZIP_DEFLATED) as zipf: folder_index = 0 while True: folder_path = self.output_file_template.format(folder_index) print(f"Checking folder: {folder_path}") if not os.path.exists(folder_path) or not os.path.isdir(folder_path): print(f"Folder {folder_path} does not exist or is not a directory. Stopping.") break # Stop if the folder doesn't exist for filename in os.listdir(folder_path): file_path = os.path.join(folder_path, filename) if filename.endswith(".csv"): print(f"Adding {file_path} to zip") zipf.write(file_path, arcname=os.path.join(os.path.basename(folder_path), filename)) else: print(f"Skipping non-CSV file: {filename}") folder_index += 1 print(f"All applicable CSV files zipped into {self.zip_file_path}")
解决方案
核心原因
Python标准库的os模块仅支持本地文件系统操作,无法识别ofs://这类分布式文件系统路径,导致os.path.exists()、os.listdir()等方法返回错误结果。必须更换为支持OFS的文件操作库或使用Spark原生的文件系统API。
可行方案
方案1:使用Spark原生FileSystem API
利用Spark的Hadoop FileSystem API直接操作OFS路径,无需额外安装依赖:
from py4j.java_gateway import java_import from pyspark.sql import SparkSession import tempfile import zipfile import os class DataFrameChunkSaver: # 保留原有的__init__、count_chunks、repartition_df、save_chunks方法不变 def zip_csv_files(self): print(f"Zipping CSV files to {self.zip_file_path}") spark = SparkSession.getActiveSession() java_import(spark._jvm, 'org.apache.hadoop.fs.Path') fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(spark._jsc.hadoopConfiguration()) # 创建本地临时zip文件,避免直接在OFS上操作大文件 with tempfile.NamedTemporaryFile(delete=False, suffix='.zip') as temp_zip: temp_zip_path = temp_zip.name with zipfile.ZipFile(temp_zip_path, 'w', zipfile.ZIP_DEFLATED) as zipf: folder_index = 0 while True: folder_path = self.output_file_template.format(folder_index) print(f"Checking folder: {folder_path}") hadoop_path = spark._jvm.Path(folder_path) if not fs.exists(hadoop_path) or not fs.isDirectory(hadoop_path): print(f"Folder {folder_path} does not exist or is not a directory. Stopping.") break # 遍历目录下的所有文件 status_list = fs.listStatus(hadoop_path) for status in status_list: file_path = status.getPath().toString() filename = status.getPath().getName() if filename.endswith(".csv"): print(f"Adding {file_path} to zip") # 将OFS文件下载到本地临时文件后添加到zip with tempfile.NamedTemporaryFile(delete=True) as temp_file: fs.copyToLocalFile(hadoop_path, spark._jvm.Path(temp_file.name)) zipf.write(temp_file.name, arcname=os.path.join(os.path.basename(folder_path), filename)) else: print(f"Skipping non-CSV file: {filename}") folder_index += 1 # 将本地zip文件上传到OFS目标路径 fs.copyFromLocalFile(spark._jvm.Path(temp_zip_path), spark._jvm.Path(self.zip_file_path)) # 清理本地临时文件 os.remove(temp_zip_path) print(f"All applicable CSV files zipped into {self.zip_file_path}")
方案2:使用hdfs3库操作OFS
如果环境允许安装第三方库,hdfs3可以更简洁地操作OFS:
- 安装依赖:
pip install hdfs3
- 修改
zip_csv_files方法:
from hdfs3 import HDFileSystem import tempfile import zipfile import os class DataFrameChunkSaver: # 保留原有的其他方法不变 def zip_csv_files(self): print(f"Zipping CSV files to {self.zip_file_path}") # 根据实际OFS集群配置初始化客户端 hdfs = HDFileSystem(host='your-ofs-host', port=your-ofs-port) # 创建本地临时zip文件 with tempfile.NamedTemporaryFile(delete=False, suffix='.zip') as temp_zip: temp_zip_path = temp_zip.name with zipfile.ZipFile(temp_zip_path, 'w', zipfile.ZIP_DEFLATED) as zipf: folder_index = 0 while True: folder_path = self.output_file_template.format(folder_index) print(f"Checking folder: {folder_path}") if not hdfs.exists(folder_path) or not hdfs.isdir(folder_path): print(f"Folder {folder_path} does not exist or is not a directory. Stopping.") break # 遍历目录下的文件 for file_full_path in hdfs.ls(folder_path): filename = os.path.basename(file_full_path) if filename.endswith(".csv"): print(f"Adding {file_full_path} to zip") # 读取OFS文件内容并写入zip with hdfs.open(file_full_path, 'rb') as f_in: with tempfile.NamedTemporaryFile(delete=True) as temp_file: temp_file.write(f_in.read()) temp_file.flush() zipf.write(temp_file.name, arcname=os.path.join(os.path.basename(folder_path), filename)) else: print(f"Skipping non-CSV file: {filename}") folder_index += 1 # 上传zip到OFS hdfs.put(temp_zip_path, self.zip_file_path) # 清理本地临时文件 os.remove(temp_zip_path) print(f"All applicable CSV files zipped into {self.zip_file_path}")
注意事项
- 确保Spark集群已配置正确的OFS访问权限,避免权限拒绝错误
- 处理超大文件时,建议采用流式读写,避免内存占用过高
- 若OFS支持直接流式写入压缩文件,可进一步优化流程,省略本地临时文件步骤
内容的提问来源于stack exchange,提问作者johnnn
相关产品推荐
相关产品推荐

