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

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中如何实现相同的任意日期范围读取逻辑?

目前已实现的内容

  1. 针对每周运行场景,生成7个单独的日期路径
  2. 针对每月运行场景,生成通配符模式: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]
    }
}

(注:修正了原代码中变量名不一致、返回类型不匹配的小问题)

待解决问题

  1. 如何生成任意日期范围的S3路径?例如:2021年11月3日至2022年3月15日
  2. 如果使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 06:27:29