Spark 不支持直接从 zip 读取/写入,因此使用 ZipOutputStream 基本上是唯一的方法。
这是我用来通过 spark 压缩现有数据的代码。它递归地列出文件的目录,然后继续压缩它们。此代码不保留目录结构,但保留文件名。
输入目录:
unzipped/
├── part-00001
├── part-00002
└── part-00003
0 directories, 3 files
输出目录:
zipped/
├── part-00001.zip
├── part-00002.zip
└── part-00003.zip
0 directories, 3 files
ZipPacker.scala:
package com.haodemon.spark.compression
import org.apache.hadoop.conf.Configuration
import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.hadoop.io.IOUtils
import org.apache.spark.sql.SparkSession
import org.apache.spark.{SparkConf, SparkContext}
import java.io.FileOutputStream
import java.util.zip.{ZipEntry, ZipOutputStream}
object ZipPacker extends Serializable {
private def getSparkContext: SparkContext = {
val conf: SparkConf = new SparkConf()
.setAppName("local")
.setMaster("local[*]")
SparkSession.builder().config(conf).getOrCreate().sparkContext
}
// recursively list files in a filesystem
private def listFiles(fs: FileSystem, path: Path): List[Path] = {
fs.listStatus(path).flatMap(p =>
if (p.isDirectory) listFiles(fs, p.getPath)
else List(p.getPath)
).toList
}
// zip compress file one by one in parallel
private def zip(inputPath: Path, outputDirectory: Path): Unit = {
val outputPath = {
val name = inputPath.getName + ".zip"
outputDirectory + "/" + name
}
println(s"Zipping to $outputPath")
val zipStream = {
val out = new FileOutputStream(outputPath)
val zip = new ZipOutputStream(out)
val entry = new ZipEntry(inputPath.getName)
zip.putNextEntry(entry)
// max compression
zip.setLevel(9)
zip
}
val conf = new Configuration
val uncompressedStream = inputPath.getFileSystem(conf).open(inputPath)
val close = true
IOUtils.copyBytes(uncompressedStream, zipStream, conf, close)
}
def main(args: Array[String]): Unit = {
val input = new Path(args(0))
println(s"Using input path $input")
val sc = getSparkContext
val uncompressedFiles = {
val conf = sc.hadoopConfiguration
val fs = input.getFileSystem(conf)
listFiles(fs, input)
}
val rdd = sc.parallelize(uncompressedFiles)
val output = new Path(args(1))
println(s"Using output path $output")
rdd.foreach(unzipped => zip(unzipped, output))
}
}