Spark Scala如何读取S3中指定日期范围的多分区数据?
问题背景
我有一个按年、月、日分区的S3数据集,目录结构如下:
year=2022/month=06/day=1/ year=2022/month=06/day=2/ year=2022/month=06/day=3/ ... year=2022/month=07/day=09/ year=2022/month=07/day=10/ year=2022/month=07/day=11/
我希望基于任意给定的日期范围,将S3存储桶中的多个分区读取到Spark Dataframe中,不想逐个读取日期分区再进行union操作。我了解PySpark中使用dateutil.relativedelta的解决方案,但想知道在Scala中如何实现相同的任意日期范围读取逻辑?
目前已实现的内容
- 针对每周运行场景,生成7个单独的日期路径
- 针对每月运行场景,生成通配符模式:
yy/MM/*
实现代码
import org.apache.spark.sql.SparkSession import java.time.{LocalDate, DayOfWeek, TemporalAdjusters} import java.time.format.DateTimeFormatter def calculate_query_dates(spark:SparkSession, runDate: LocalDate, programType:String): (LocalDate, LocalDate) = { val dayOfWeek = runDate.getDayOfWeek.getValue() if (programType == "Month") { val last_day_prev_month = runDate.minusMonths(1).`with`(TemporalAdjusters.lastDayOfMonth()) val first_day_prev_month = runDate.minusMonths(1).`with`(TemporalAdjusters.firstDayOfMonth()) (first_day_prev_month, last_day_prev_month) } else if (programType == "Week"){ val last_day_prev_week = if (dayOfWeek == 7) runDate else runDate.`with`(TemporalAdjusters.previous( DayOfWeek.SATURDAY)) println("last_day_of_week:::", last_day_prev_week) val first_day_prev_week = if (dayOfWeek == 7) runDate.`with`(TemporalAdjusters.previous( DayOfWeek.SUNDAY)) else runDate.minusDays(7).`with`(TemporalAdjusters.previous( DayOfWeek.SUNDAY)) println("first_day_of_week:::", first_day_prev_week) (first_day_prev_week,last_day_prev_week) } else { // 补充默认分支避免编译报错 (runDate, runDate) } } def generate_S3_paths(spark:SparkSession, s3_bucket_name: String, program_type: String, s3_prefix: String ="", start_query_date: LocalDate, end_query_date: LocalDate): List[String] = { val monthFormatter = DateTimeFormatter.ofPattern("MM"); val dateFormatter = DateTimeFormatter.ofPattern("dd"); if (program_type == "Month"){ List(s3_bucket_name + "/" + s3_prefix + "year=" + start_query_date.getYear.toString + "/month=" + start_query_date.format(monthFormatter).toString + "/*") } else if (program_type == "Week"){ val list_dates = (0 until 7).map(start_query_date.plusDays(_)).toList list_dates.map(i => s3_bucket_name+"/"+ s3_prefix + "year=" + i.getYear.toString + "/month=" + i.format(monthFormatter).toString + "/day=" + i.format(dateFormatter).toString + "/") } else { List.empty[String] } }
(注:修正了原代码中变量名不一致、返回类型不匹配的小问题)
待解决问题
- 如何生成任意日期范围的S3路径?例如:2021年11月3日至2022年3月15日
- 如果使用
parquet(paths: String*)读取这些S3路径(如list_s3_paths_week),是顺序读取还是并行读取?
解答
1. 生成任意日期范围的S3路径
Scala中可以通过遍历日期范围生成路径,同时优化逻辑:对完整月份直接用day=*通配符,减少路径数量,提升Spark扫描效率;对非完整月份则生成单个日期的路径。
实现代码
import java.time.{LocalDate, TemporalAdjusters} import java.time.format.DateTimeFormatter def generateArbitraryDateRangePaths( s3Bucket: String, s3Prefix: String, startDate: LocalDate, endDate: LocalDate ): List[String] = { val monthFormatter = DateTimeFormatter.ofPattern("MM") val dayFormatter = DateTimeFormatter.ofPattern("dd") val basePath = s"$s3Bucket/$s3Prefix" // 获取日期范围覆盖的所有月份 val currentMonth = startDate.with(TemporalAdjusters.firstDayOfMonth()) val endMonth = endDate.with(TemporalAdjusters.firstDayOfMonth()) val months = Iterator.iterate(currentMonth)(_.plusMonths(1)) .takeWhile(!_.isAfter(endMonth)) .toList months.flatMap { month => val isStartMonth = month.isEqual(currentMonth) val isEndMonth = month.isEqual(endMonth) if (isStartMonth && isEndMonth) { // 起止在同一个月,生成该月内的所有日期路径 val days = Iterator.iterate(startDate)(_.plusDays(1)) .takeWhile(!_.isAfter(endDate)) .toList days.map { day => s"$basePath/year=${day.getYear}/month=${day.format(monthFormatter)}/day=${day.format(dayFormatter)}/" } } else if (isStartMonth) { // 起始月,生成从startDate到月末的路径 val endOfMonth = month.with(TemporalAdjusters.lastDayOfMonth()) val days = Iterator.iterate(startDate)(_.plusDays(1)) .takeWhile(!_.isAfter(endOfMonth)) .toList days.map { day => s"$basePath/year=${day.getYear}/month=${day.format(monthFormatter)}/day=${day.format(dayFormatter)}/" } } else if (isEndMonth) { // 结束月,生成从月初到endDate的路径 val startOfMonth = month.with(TemporalAdjusters.firstDayOfMonth()) val days = Iterator.iterate(startOfMonth)(_.plusDays(1)) .takeWhile(!_.isAfter(endDate)) .toList days.map { day => s"$basePath/year=${day.getYear}/month=${day.format(monthFormatter)}/day=${day.format(dayFormatter)}/" } } else { // 中间完整月,直接用通配符 List(s"$basePath/year=${month.getYear}/month=${month.format(monthFormatter)}/day=*/") } } }
示例使用
val start = LocalDate.of(2021, 11, 3) val end = LocalDate.of(2022, 3, 15) val paths = generateArbitraryDateRangePaths("s3://my-bucket", "data/", start, end) // 生成的路径包含: // 2021-11-03至2021-11-30的单个日期路径 // 2021-12、2022-01、2022-02的通配符路径 // 2022-03-01至2022-03-15的单个日期路径
2. Spark读取多路径的并行性
使用parquet(paths: String*)读取多个路径时,Spark会并行读取:
- Spark会将每个路径对应的文件/分区作为独立任务,分配到集群不同Executor上执行,无需手动union——
parquet方法本身支持合并多路径为单个DataFrame,性能远优于逐个读取再union。 - 通配符路径的效率更高,因为Spark可直接通过目录结构扫描符合条件的文件,减少客户端生成的路径数量,降低调度开销。
内容的提问来源于stack exchange,提问作者riyaB
相关产品推荐
相关产品推荐

