Scala 2.11.7获取tar压缩包内csv.gz文件的实现方法
Hey,我刚好处理过类似的需求,针对Scala 2.11.7下从指定目录的所有tar包中提取.csv.gz文件列表,再转为DataFrame做后续处理的场景,咱们一步步来实现:
1. 先准备依赖
要处理tar压缩包和gzip文件,咱们需要用到Apache Commons Compress库,它对Scala 2.11.7兼容性很好。如果用SBT管理依赖,在build.sbt里添加:
libraryDependencies += "org.apache.commons" % "commons-compress" % "1.20"
(这个版本是经过验证兼容Scala 2.11的,你也可以根据需求调整,但尽量选稳定版)
2. 遍历目录下的所有.tar文件
首先写个工具方法,获取目标目录里所有后缀为.tar的文件:
import java.io.File def listTarFiles(dirPath: String): List[File] = { val dir = new File(dirPath) if (dir.isDirectory) { dir.listFiles().filter(_.getName.endsWith(".tar")).toList } else { Nil // 如果传入的不是目录,返回空列表 } }
3. 从单个Tar包中筛选.csv.gz文件
接下来,针对每个tar包,用TarArchiveInputStream读取里面的条目,筛选出以.csv.gz结尾的文件(注意排除目录条目):
import org.apache.commons.compress.archivers.tar.TarArchiveInputStream import java.io.FileInputStream def listCsvGzInTar(tarFile: File): List[String] = { val tarIn = new TarArchiveInputStream(new FileInputStream(tarFile)) var entry = tarIn.getNextTarEntry() val csvGzEntries = scala.collection.mutable.ListBuffer[String]() while (entry != null) { // 只保留非目录且后缀为.csv.gz的条目 if (!entry.isDirectory && entry.getName.endsWith(".csv.gz")) { csvGzEntries += entry.getName } entry = tarIn.getNextTarEntry() } tarIn.close() // 记得关闭流,避免资源泄漏 csvGzEntries.toList }
4. 整合获取所有.csv.gz文件列表
把上面两个方法结合起来,就能拿到目录下所有tar包中的.csv.gz文件列表了,我们还可以把tar文件和对应的csv.gz路径配对,方便后续处理:
def getAllCsvGzPaths(tarDir: String): List[(File, String)] = { listTarFiles(tarDir).flatMap { tarFile => listCsvGzInTar(tarFile).map(csvGzPath => (tarFile, csvGzPath)) } } // 调用示例:替换成你的tar目录路径 val tarDirectory = "/path/to/your/tar/files" val allCsvGz = getAllCsvGzPaths(tarDirectory) // 打印验证结果 allCsvGz.foreach { case (tarFile, csvGzPath) => println(s"找到文件:${csvGzPath},所在Tar包:${tarFile.getName}") }
5. 将.csv.gz文件转为DataFrame
接下来是核心的转DataFrame步骤。因为文件在tar包里,不能直接用Spark的常规read.csv,需要读取tar包内的gzip流,再转换成Spark能处理的数据源。这里假设你用的是Spark 2.x(对应Scala 2.11):
首先初始化SparkSession:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.types._ val spark = SparkSession.builder() .appName("TarCsvGzToDataFrame") .master("local[*]") // 生产环境请移除这个配置,用集群模式 .getOrCreate() import spark.implicits._
然后定义你的CSV文件的Schema(根据实际数据结构调整):
val csvSchema = StructType(Array( StructField("id", StringType, nullable = true), StructField("value", DoubleType, nullable = true), StructField("timestamp", TimestampType, nullable = true) ))
接下来写方法处理单个.csv.gz条目,转换成DataFrame:
import org.apache.commons.compress.compressors.gzip.GzipCompressorInputStream import scala.io.Source def convertCsvGzToDataFrame(tarFile: File, csvGzPath: String) = { try { val tarIn = new TarArchiveInputStream(new FileInputStream(tarFile)) var entry = tarIn.getNextTarEntry() // 定位到目标.csv.gz条目 while (entry != null && entry.getName != csvGzPath) { entry = tarIn.getNextTarEntry() } entry match { case null => println(s"在${tarFile.getName}中未找到${csvGzPath}") spark.emptyDataFrame case _ => try { // 读取gzip压缩流 val gzIn = new GzipCompressorInputStream(tarIn) // 把流转换成字符串行 val csvLines = Source.fromInputStream(gzIn).getLines().toList // 并行化后转DataFrame val csvRdd = spark.sparkContext.parallelize(csvLines) spark.read.schema(csvSchema).csv(csvRdd) } finally { tarIn.close() // 自动关闭关联的gzIn } } } catch { case e: Exception => println(s"处理${csvGzPath}失败:${e.getMessage}") spark.emptyDataFrame } }
最后批量处理所有.csv.gz文件,合并成一个大的DataFrame(如果所有文件结构一致):
// 批量转换所有文件 val allDataFrames = allCsvGz.map { case (tar, csvGz) => convertCsvGzToDataFrame(tar, csvGz) } // 合并所有DataFrame(结构一致时可用) val combinedDf = allDataFrames.reduce(_ unionByName _) // 后续数据转换操作示例 val transformedDf = combinedDf .filter($"value" > 0) .withColumn("date", to_date($"timestamp")) transformedDf.show()
注意事项
- 资源泄漏:一定要确保流被正确关闭,上面的代码用了
try...finally块来保证,你也可以用Scala的Using语法(如果是Scala 2.13+,但2.11没有,所以还是用try-finally更稳妥)。 - 性能优化:如果tar包很大,建议不要一次性把所有行读到内存里,可以考虑用Spark的自定义InputFormat来直接读取tar包内的文件,避免内存溢出。
- Schema适配:一定要根据你的实际CSV结构调整Schema,否则会出现数据类型不匹配的问题。
内容的提问来源于stack exchange,提问作者Chaouki
相关产品推荐
相关产品推荐

