如何用Scala编写Glue作业读取S3存储桶中最新文件夹的数据?
Scala实现Glue作业读取S3最新文件夹数据
实现方案
核心逻辑是先枚举S3桶内的目标文件夹,解析文件夹名称的时间戳并筛选出最新路径,最后读取该路径下的数据。以下是具体实现步骤与代码:
1. 前置准备
确保Glue作业的IAM角色拥有s3:ListBucket(枚举文件夹)和s3:GetObject(读取数据)权限,同时替换代码中的桶名、基础路径等配置项。
2. 完整Scala代码示例
import com.amazonaws.services.s3.AmazonS3ClientBuilder import com.amazonaws.services.s3.model.ListObjectsV2Request import java.time.LocalDateTime import java.time.format.DateTimeFormatter import com.amazonaws.services.glue.GlueContext import org.apache.spark.SparkContext import scala.collection.JavaConverters._ object LatestS3FolderGlueJob { def main(args: Array[String]): Unit = { // 初始化Glue上下文 val sc: SparkContext = new SparkContext() val glueContext: GlueContext = new GlueContext(sc) val spark = glueContext.getSparkSession // 配置S3参数 val bucketName = "your-target-bucket" val basePrefix = "your-base-folder/" // 无基础路径则设为"" // 初始化S3客户端 val s3Client = AmazonS3ClientBuilder.defaultClient() // 构造请求,按"/"分隔枚举所有前缀(模拟文件夹) val listRequest = new ListObjectsV2Request() .withBucketName(bucketName) .withPrefix(basePrefix) .withDelimiter("/") // 获取所有文件夹前缀 val s3Result = s3Client.listObjectsV2(listRequest) val prefixes = s3Result.getCommonPrefixes.asScala // 定义时间格式,匹配文件夹名称"YYYY-MM-DD HH:MM:SS" val dateFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss") // 筛选并找出最新文件夹路径 val latestFolder = prefixes .map(_.replace(basePrefix, "").stripSuffix("/")) // 提取纯文件夹名称 .filter(folderName => { // 校验时间格式合法性 try { LocalDateTime.parse(folderName, dateFormatter) true } catch { case _: Exception => false } }) .map(folderName => (LocalDateTime.parse(folderName, dateFormatter), folderName)) .maxByOption(_._1) // 取时间戳最大的文件夹 .map { case (_, name) => s"s3://$bucketName/$basePrefix$name/" } // 读取最新文件夹数据 latestFolder match { case Some(path) => println(s"读取目标路径: $path") // 用Glue DynamicFrame读取(支持自动推断Schema,可根据数据格式调整) val dynamicFrame = glueContext.getSourceWithFormat( connectionType = "s3", format = "parquet", // 替换为实际数据格式:json/csv等 formatOptions = JsonOptions(Map("inferSchema" -> "true")), options = JsonOptions(Map("path" -> path)) ).getDynamicFrame() // 转换为Spark DataFrame进行后续业务处理 val df = dynamicFrame.toDF() df.show(10) // 此处添加自定义业务逻辑 case None => println("未找到符合格式的有效文件夹") } } }
3. 关键说明
- 时间格式匹配:严格按照文件夹名称的
YYYY-MM-DD HH:MM:SS格式解析,非法格式的文件夹会被自动过滤。 - 分页处理:如果S3桶内文件夹数量极大,需添加循环调用
s3Client.listNextBatchOfObjects()的逻辑,处理分页返回的结果。 - 数据格式适配:根据实际存储的数据类型修改
format参数(如csv需额外配置separator等formatOptions)。
内容的提问来源于stack exchange,提问作者NikRED
相关产品推荐
相关产品推荐

