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

如何使用PySpark或Spark Scala对FTP服务器zip文件进行解压及文件迁移

实现方案说明

你使用的com.springml.spark.sftp第三方包仅支持直接读取SFTP上的非压缩结构化文件(如CSV、JSON),不支持zip压缩包的解压、文件筛选、文件回传FTP的操作,且原有代码未指定具体读取路径,所以无法正常运行。

Scala 实现步骤

前置依赖

需要在项目中引入以下依赖(以sbt为例):

libraryDependencies ++= Seq(
  "commons-net" % "commons-net" % "3.9.0", // FTP操作工具
  "org.apache.commons" % "commons-compress" % "1.23.0" // 压缩包处理工具
)

完整实现代码

import org.apache.commons.net.ftp.FTPClient
import org.apache.commons.compress.archivers.zip.ZipFile
import java.io.{File, FileOutputStream, InputStream}
import scala.collection.JavaConverters._

// 自定义配置项
val FTP_HOST = "你的FTP地址"
val FTP_PORT = 21
val FTP_USER = "用户名"
val FTP_PWD = "密码"
val ZIP_REMOTE_PATH = "/ftp源目录/目标压缩包.zip"
val LOCAL_TEMP_DIR = "/临时存储路径/" // 集群环境建议用HDFS等分布式存储路径
val TARGET_FILE_SUFFIX = ".csv" // 自定义筛选条件,可替换为文件名前缀、全名匹配等规则
val FTP_TARGET_DIR = "/ftp目标目录/"

// 1. 从FTP拉取zip文件到临时目录
val ftpClient = new FTPClient()
ftpClient.connect(FTP_HOST, FTP_PORT)
ftpClient.login(FTP_USER, FTP_PWD)
ftpClient.enterLocalPassiveMode()

val localZipFile = new File(LOCAL_TEMP_DIR + "temp.zip")
val outputStream = new FileOutputStream(localZipFile)
ftpClient.retrieveFile(ZIP_REMOTE_PATH, outputStream)
outputStream.close()

// 2. 解压zip并筛选目标文件
val zipFile = new ZipFile(localZipFile)
val targetFiles = zipFile.getEntries.asScala.filter(entry => 
  !entry.isDirectory && entry.getName.endsWith(TARGET_FILE_SUFFIX)
).toList

// 3. 将筛选后的文件上传到FTP目标目录
targetFiles.foreach(entry => {
  val inputStream: InputStream = zipFile.getInputStream(entry)
  val remoteFileName = FTP_TARGET_DIR + entry.getName.split("/").last
  ftpClient.storeFile(remoteFileName, inputStream)
  inputStream.close()
})

// 4. 可选:如果需要读取目标文件内容为DataFrame做后续处理
val df = spark.read
  .option("delimiter", ";")
  .option("quote", "\"")
  .option("escape", "\\")
  .option("multiLine", "true")
  .option("inferSchema", "true")
  .csv(LOCAL_TEMP_DIR + "解压后的文件路径")

// 资源释放
zipFile.close()
ftpClient.logout()
ftpClient.disconnect()
localZipFile.delete() // 删除临时文件
PySpark 实现步骤

前置依赖

安装需要的Python包:

pip install ftplib zipfile36

完整实现代码

from ftplib import FTP
import zipfile
import os
from pyspark.sql import SparkSession

# 自定义配置项
FTP_HOST = "你的FTP地址"
FTP_PORT = 21
FTP_USER = "用户名"
FTP_PWD = "密码"
ZIP_REMOTE_PATH = "/ftp源目录/目标压缩包.zip"
LOCAL_TEMP_DIR = "/临时存储路径/"
TARGET_FILE_SUFFIX = ".csv"
FTP_TARGET_DIR = "/ftp目标目录/"

spark = SparkSession.builder.appName("FtpZipProcess").getOrCreate()

# 1. 拉取zip到临时目录
ftp = FTP()
ftp.connect(FTP_HOST, FTP_PORT)
ftp.login(FTP_USER, FTP_PWD)
ftp.set_pasv(True)

local_zip_path = os.path.join(LOCAL_TEMP_DIR, "temp.zip")
with open(local_zip_path, 'wb') as f:
    ftp.retrbinary(f'RETR {ZIP_REMOTE_PATH}', f.write)

# 2. 解压筛选文件
target_files = []
with zipfile.ZipFile(local_zip_path, 'r') as zf:
    for file_info in zf.infolist():
        if not file_info.is_dir() and file_info.filename.endswith(TARGET_FILE_SUFFIX):
            zf.extract(file_info, LOCAL_TEMP_DIR)
            target_files.append(os.path.join(LOCAL_TEMP_DIR, file_info.filename))

# 3. 上传到FTP目标目录
for file_path in target_files:
    file_name = os.path.basename(file_path)
    with open(file_path, 'rb') as f:
        ftp.storbinary(f'STOR {os.path.join(FTP_TARGET_DIR, file_name)}', f)
    os.remove(file_path)

# 4. 可选:读取目标文件为DataFrame
df = spark.read.csv(
    target_files,
    sep=";",
    quote='"',
    escape="\\",
    multiLine=True,
    inferSchema=True
)

# 资源释放
os.remove(local_zip_path)
ftp.quit()
注意事项
  • 如果是SFTP协议,Scala侧将FTPClient替换为JSch工具类,Python侧将ftplib替换为paramiko库即可,整体处理逻辑基本一致
  • 集群运行时,请将临时存储路径替换为HDFS、对象存储等分布式共享存储路径,避免单节点本地磁盘无法被其他节点访问的问题
  • 体积超过1G的大压缩包建议使用spark.sparkContext.binaryFiles读取zip二进制流,做分布式解压处理,避免单节点CPU、内存压力过大
  • 操作完成后请及时删除临时压缩包和解压后的中间文件,避免占用存储资源

内容的提问来源于stack exchange,提问作者Omayma HARBAOUI

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 11:30:05