使用PySpark从SFTP读取大Parquet文件速度过慢的问题
问题描述
我使用SQLContext读取SFTP服务器上的Parquet文件时遇到性能瓶颈,该文件包含600万行数据,现有方案读取耗时近1小时。
当前可运行但速度极慢的脚本:
import pyarrow as pa import pyarrow.parquet as pq from fsspec.implementations.sftp import SFTPFileSystem fs = SFTPFileSystem(host = SERVER_SFTP, port = SERVER_PORT, username = USER, password = PWD) df = pq.read_table(SERVER_LOCATION\FILE.parquet, filesystem = fs)
本地读取大文件时,以下代码效率很高,因此我想了解如何用SparkSQL直接读取SFTP上的远程文件:
df = sqlContext.read.parquet('PATH/file')
我尝试过两种方法但均未达到预期:
- 直接用SFTP库打开文件传给Spark,完全丧失了Spark的并行处理优势:
df = sqlContext.read.parquet(sftp.open('PATH/file')) - 尝试使用spark-sftp库,但未成功。
解决方案
方法1:利用Spark的Hadoop SFTP文件系统支持
Spark底层依赖Hadoop生态,可通过配置Hadoop的SFTP文件系统实现Spark直接读取SFTP文件,步骤如下:
- 确认依赖:确保Spark作业环境包含
hadoop-common和hadoop-sftp相关依赖(多数Hadoop发行版自带,若缺失需手动引入对应版本的jar包)。 - 配置SparkSession:在初始化时传入SFTP的认证与连接参数,之后即可像读取本地文件一样读取SFTP路径:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("SFTPParquetReader") \ .config("spark.hadoop.fs.sftp.impl", "org.apache.hadoop.fs.sftp.SFTPFileSystem") \ .config("spark.hadoop.fs.sftp.host", SERVER_SFTP) \ .config("spark.hadoop.fs.sftp.port", SERVER_PORT) \ .config("spark.hadoop.fs.sftp.user", USER) \ .config("spark.hadoop.fs.sftp.password", PWD) \ .getOrCreate() sqlContext = spark.sqlContext # 注意路径格式为sftp://开头的远程路径 df = sqlContext.read.parquet(f"sftp://{SERVER_LOCATION}/FILE.parquet")
方法2:并行化下载后读取(折中方案)
若Hadoop配置存在困难,可借助Spark的并行能力分片下载SFTP上的Parquet文件(适用于分块存储的Parquet),再读取本地临时文件:
from fsspec.implementations.sftp import SFTPFileSystem from pyspark.sql import SparkSession spark = SparkSession.builder.appName("SFTPParallelDownload").getOrCreate() fs = SFTPFileSystem(host=SERVER_SFTP, port=SERVER_PORT, username=USER, password=PWD) # 列出远程Parquet的分块文件(假设是按part-*命名的分块) remote_part_files = fs.glob(f"{SERVER_LOCATION}/FILE.parquet/part-*.parquet") # 定义并行下载函数 def download_to_local(remote_path): local_temp_path = f"/tmp/sftp_parquet_temp/{remote_path.split('/')[-1]}" fs.download(remote_path, local_temp_path) return local_temp_path # 用Spark RDD并行执行下载任务 spark.sparkContext.parallelize(remote_part_files).map(download_to_local).collect() # 读取本地临时目录的Parquet文件,自动利用Spark并行优势 df = sqlContext.read.parquet("/tmp/sftp_parquet_temp/")
方法3:优化Parquet文件结构提升并行度
如果远程Parquet是单个大文件,会限制Spark的并行处理能力。建议:
- 上传到SFTP前,先将Parquet拆分为多个分块文件(比如通过Spark写入时指定分区数)。
- 若无法修改源文件,读取后可通过
repartition重新分区,后续处理能利用并行优势:df = sqlContext.read.parquet(f"sftp://{SERVER_LOCATION}/FILE.parquet").repartition(10)
内容的提问来源于stack exchange,提问作者A2N15
相关产品推荐
相关产品推荐

