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

