如何使用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
相关产品推荐
相关产品推荐

