You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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:

  1. 安装依赖:
pip install hdfs3
  1. 修改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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 16:09:56