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

Spark Scala脚本每日读取S3按日期存储CSV文件的实现方案

解决方案:动态生成S3日期路径

针对你的需求,这里有几个实用的方案,帮你在Spark Scala脚本里自动匹配S3上的日期文件夹:

1. 直接在Scala中生成当前系统日期(最常用)

利用Java 8+的时间API(线程安全,推荐替代旧的SimpleDateFormat)生成指定格式的日期字符串,再拼接成S3路径:

import java.time.LocalDate
import java.time.format.DateTimeFormatter
import org.apache.spark.SparkContext

// 初始化SparkContext(如果还没初始化)
val sc = new SparkContext()

// 定义和S3文件夹一致的日期格式:MM-DD-YYYY
val dateFormatter = DateTimeFormatter.ofPattern("MM-dd-yyyy")
// 获取当前系统日期并格式化
val todayDate = LocalDate.now().format(dateFormatter)

// 动态拼接S3路径
val s3Path = s"s3a://digital/$todayDate/abc.csv"
val fileFromS3 = sc.textFile(s3Path)

如果需要指定时区(比如S3文件夹用UTC时间,而你的Spark集群是其他时区),可以加上时区参数:

import java.time.ZoneId
val todayDate = LocalDate.now(ZoneId.of("UTC")).format(dateFormatter)

2. 通过命令行参数传递日期(适合回溯任务)

如果需要偶尔跑历史日期的任务,而不是固定读取当天,可以在提交Spark作业时传入日期参数,脚本里读取参数生成路径:

提交作业时传参:

spark-submit \
  --class com.your.package.YourMainClass \
  --master yarn \
  your-spark-app.jar \
  09-25-2024

脚本中读取参数:

import org.apache.spark.SparkContext

val sc = new SparkContext()
// 获取命令行传入的日期参数
val targetDate = sc.getConf.get("spark.driver.args").split(" ")(0)

val s3Path = s"s3a://digital/$targetDate/abc.csv"
val fileFromS3 = sc.textFile(s3Path)

你也可以用--conf参数更清晰地传递:

spark-submit \
  --class com.your.package.YourMainClass \
  --master yarn \
  --conf "spark.target.date=09-25-2024" \
  your-spark-app.jar

脚本中读取:

val targetDate = sc.getConf.get("spark.target.date")

3. 结合调度工具生成日期(生产环境最佳实践)

如果你的Spark作业是通过调度工具(比如Airflow、Oozie)自动触发的,可以直接在调度层生成日期字符串,再传给Spark脚本。比如Airflow中可以用日期宏:

  • 当天日期:{{ execution_date.strftime('%m-%d-%Y') }}
  • 自定义格式完全匹配S3文件夹的格式,然后把这个值作为参数传给spark-submit,脚本逻辑和方案2一致。

这种方式的好处是把日期管理交给调度工具,更便于回溯、重试和维护作业。


内容的提问来源于stack exchange,提问作者Ganesh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 07:36:41