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
相关产品推荐
相关产品推荐

