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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:25:46