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

如何使用Spark应用实现定时删除S3中超过指定创建时长的文件

实现方案

1 逻辑前置说明

你需求中描述的文件创建时间大于当前时间戳减去1小时的差值大概率是笔误:如果当前时间是10点,当前时间减1小时是9点,创建时间大于9点的是近1小时内的新文件,通常清理逻辑是删除创建时间小于9点的1小时前的旧文件,以下实现默认按这个逻辑开发,如果确实要删除近1小时的新文件,修改判断条件即可。

2 依赖准备

需要在Spark项目中引入对应Hadoop版本的AWS依赖,Maven坐标参考如下:

<dependency>
    <groupId>org.apache.hadoop</groupId>
    <artifactId>hadoop-aws</artifactId>
    <version>和你集群Hadoop版本一致</version>
</dependency>

3 完整代码实现(Scala版本)

import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.spark.sql.SparkSession
import java.util.concurrent.Executors
import java.util.concurrent.TimeUnit

object S3HourlyCleaner {
  def main(args: Array[String]): Unit = {
    // 初始化SparkSession
    val spark = SparkSession.builder()
      .appName("S3HourlyFileCleaner")
      // 本地测试时放开下一行
      // .master("local[*]")
      .getOrCreate()

    // S3权限配置,优先用集群IAM角色/环境变量注入,不要硬编码密钥
    val hadoopConf = spark.sparkContext.hadoopConfiguration
    hadoopConf.set("fs.s3a.access.key", "你的S3访问密钥AK")
    hadoopConf.set("fs.s3a.secret.key", "你的S3密钥SK")
    hadoopConf.set("fs.s3a.endpoint", "对应S3区域的服务端点")

    // 基础参数配置
    val targetS3Path = "s3a://你的桶名/要清理的路径前缀/"
    val expireTimeMs = 3600 * 1000 // 1小时的毫秒数

    // 定义单次清理任务
    val cleanTask = new Runnable {
      override def run(): Unit = {
        val fs = FileSystem.get(hadoopConf)
        val currentTime = System.currentTimeMillis()
        val cutoffTime = currentTime - expireTimeMs

        // 递归列出目标路径下所有文件
        val allFiles = fs.listStatus(new Path(targetS3Path))
        allFiles.foreach { fileStatus =>
          // 跳过目录只处理文件
          if (!fileStatus.isDirectory) {
            // S3无原生创建时间属性,getModificationTime返回的就是文件上传时间,等价于创建时间
            val fileCreateTime = fileStatus.getModificationTime
            // 按需求调整判断条件,当前逻辑是删除1小时前的旧文件
            if (fileCreateTime < cutoffTime) {
              val delResult = fs.delete(fileStatus.getPath, false)
              if (delResult) {
                println(s"删除成功:${fileStatus.getPath.toString}")
              } else {
                println(s"删除失败:${fileStatus.getPath.toString}")
              }
            }
          }
        }
      }
    }

    // 定时调度:启动立即执行第一次,之后每1小时执行一次
    val scheduler = Executors.newScheduledThreadPool(1)
    scheduler.scheduleAtFixedRate(cleanTask, 0, 1, TimeUnit.HOURS)

    // 保持进程不退出
    spark.streams.awaitAnyTermination()
  }
}

4 优化&注意事项

  • 如果要清理的文件量级非常大,可以把文件列表转成RDD并行处理,提升清理效率:
    import spark.implicits._
    val fileDf = fs.listStatus(new Path(targetS3Path)).filter(!_.isDirectory).toSeq.toDF("file")
    fileDf.foreach { row =>
      val file = row.getAs[org.apache.hadoop.fs.FileStatus]("file")
      // 后面判断删除逻辑和上面一致
    }
    
  • 测试阶段可以先把fs.delete替换成打印文件路径,确认判断逻辑符合预期后再开启真实删除,避免误删
  • 如果不需要用到Spark分布式能力,这个场景也可以直接用AWS S3 SDK加本地定时任务实现,资源消耗更低

内容的提问来源于stack exchange,提问作者Ayushi Dewangan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 23:57:04