Scala Spark中如何用通配符判断S3指定模式对象是否存在?
解决S3中匹配
part-*.json文件存在性的问题 方案1:S3前缀遍历+正则匹配
S3不支持通配符查询,但可以通过前缀定位缩小查询范围,再用正则过滤出目标文件。由于你明确S3中仅有一个符合part-*.json的文件,只需判断过滤结果是否非空即可。
Scala代码示例:
import com.amazonaws.services.s3.AmazonS3ClientBuilder import com.amazonaws.services.s3.model.ListObjectsV2Request import scala.collection.JavaConverters._ // 初始化S3客户端 val s3Client = AmazonS3ClientBuilder.defaultClient() val bucket = "xyz" val prefix = "abc/def/part-" // 前缀锁定到part-开头的路径 val partJsonRegex = """part-.*\.json$""".r // 匹配part-*.json的正则表达式 // 构造前缀查询请求 val request = new ListObjectsV2Request() .withBucketName(bucket) .withPrefix(prefix) // 执行查询并过滤目标文件 val targetObjects = s3Client.listObjectsV2(request).getObjectSummaries.asScala .map(_.getKey) .filter(key => partJsonRegex.findFirstIn(key).isDefined) // 判断文件是否存在 val fileExists = targetObjects.nonEmpty
方案2:Spark生成时直接获取文件名(更高效)
既然文件是Spark生成的,且你是通过coalesce(1)或repartition(1)生成单个part-*.json文件,可以在写入完成后直接读取输出目录的文件列表,拿到具体文件名,后续直接用该文件名调用doesObjectExist,省去S3查询开销。
Scala代码示例:
import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.spark.sql.SparkSession val spark = SparkSession.builder().getOrCreate() val outputDir = "s3a://xyz/abc/def/" val fs = FileSystem.get(spark.sparkContext.hadoopConfiguration) // 遍历输出目录,过滤出part-*.json文件 val partFiles = fs.listStatus(new Path(outputDir)) .filter(status => !status.isDirectory && status.getPath.getName.matches("part-.*\\.json")) if (partFiles.nonEmpty) { val exactKey = s"abc/def/${partFiles.head.getPath.getName}" val fileExists = s3Client.doesObjectExist("xyz", exactKey) // 后续业务逻辑 }
方案3:前缀+后缀简化过滤
如果不想用正则,也可以直接通过前缀和后缀组合判断:
val result = s3Client.listObjectsV2(bucket, prefix) val fileExists = result.getObjectSummaries.asScala.exists { summary => val key = summary.getKey key.startsWith(prefix) && key.endsWith(".json") }
内容的提问来源于stack exchange,提问作者Mikasa
相关产品推荐
相关产品推荐

