基于Scala动态获取指定小时数的HDFS分区路径列表
需求与现有问题
- 业务场景:需根据指定小时数收集HDFS上的文件路径列表,路径按
partition_date(格式yyyy-MM-dd)和hour(格式两位数字,如00、23)分区,基础路径固定为/user/hdfs/test/,需结合当前时间自动拼接分区信息。 - 示例:
- 当前时间
2022-09-21 19:00,传入小时数2时,需生成路径:/user/hdfs/test/partition_date=2022-09-21/hour=19、/user/hdfs/test/partition_date=2022-09-21/hour=18 - 当前时间
2022-09-22 00:00,传入小时数2时,需生成跨天路径:/user/hdfs/test/partition_date=2022-09-22/hour=00、/user/hdfs/test/partition_date=2022-09-21/hour=23
- 当前时间
- 现有代码局限:仅能硬编码小时数生成路径,无法支持从Spark Submit传入动态小时数,且逻辑仅处理2小时场景,不通用。
通用实现方案
1. 通过Spark Submit传递参数
使用--conf参数传递小时数,示例命令:
spark-submit \ --class your.main.Class \ --conf spark.app.hours=3 \ your-jar-file.jar
2. 代码中读取并校验参数
从SparkConf中读取传入的小时数,确保参数为合法正整数,避免非法输入。
3. 通用时间偏移计算
用Java 8+的LocalDateTime处理时间偏移,自动处理跨天、跨月场景,无需手动判断日期切换逻辑。
4. 动态生成路径列表
循环计算每个偏移小时对应的分区日期和小时,拼接成完整HDFS路径,生成路径列表。
完整Scala代码实现
import org.apache.spark.SparkConf import java.time.{LocalDateTime, DateTimeFormatter} object HdfsPathGenerator { def main(args: Array[String]): Unit = { val sparkConf = new SparkConf().setAppName("HdfsPathGenerator") // 读取传入的小时数参数,未指定则抛出异常 val hours = sparkConf.getOption("spark.app.hours") .map(_.toInt) .filter(_ > 0) .getOrElse(throw new IllegalArgumentException("必须通过--conf spark.app.hours=X指定正整数小时数")) // 定义分区格式的时间格式化器 val dateFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd") val hourFormatter = DateTimeFormatter.ofPattern("HH") // 获取当前时间 val now = LocalDateTime.now() // 生成所有目标路径 val paths = (0 until hours).map { offset => val targetTime = now.minusHours(offset) val partitionDate = targetTime.format(dateFormatter) val hour = targetTime.format(hourFormatter) s"/user/hdfs/test/partition_date=$partitionDate/hour=$hour/" } // 输出路径(可根据业务需求替换为文件读取等逻辑) println("生成的HDFS路径列表:") paths.foreach(println) // 后续业务逻辑... } }
代码说明
- 参数校验:确保传入的小时数为正整数,避免非法参数导致错误。
- 时间处理:
LocalDateTime.minusHours()自动处理时间偏移,无需手动判断跨天,逻辑简洁可靠。 - 格式规范:用
DateTimeFormatter保证日期和小时格式符合分区要求(小时为两位数字,如00而非0)。 - 扩展性:支持任意正整数小时数的需求,无需修改核心逻辑。
内容的提问来源于stack exchange,提问作者Rahul Patidar
相关产品推荐
相关产品推荐

