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

使用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文件,步骤如下:

  1. 确认依赖:确保Spark作业环境包含hadoop-common和hadoop-sftp相关依赖(多数Hadoop发行版自带,若缺失需手动引入对应版本的jar包)。
  2. 配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 18:19:47