所在公司使用的线上调度是 Azkaban,重跑间隔段很长的数据时很麻烦,申请使用dolphinscheduler重跑,但是暂时未允许,并且后来想想所在公司的代码多是spark和shell脚本,强行使用ds的话还是需要重新配置项目,配置环境,也是麻烦,就想着写个spark循环得了,记录一下工具代码吧,万一哪天跑路了也能找得到。

import java.util.Date
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions._

/** *********************
 * *******
 * Create by wangk  ********
 * 2026/1/8      *******
 * **********************/
object transferDate {

  def main(args: Array[String]): Unit = {
    // 创建SparkSession
    val spark = SparkSession.builder()
      .appName("YearDatePrinter")
      .master("local[*]")
      .getOrCreate()

    import spark.implicits._

    // 示例数据集
    val data = Seq(("2025-10-31", "2025-12-31","p"))
    val df = data.toDF( ("curr_date"), ("end_date"),("p"))//添加对应的命名
    df.show()
   /***
     和想要循环的结束时间end_date
     用来计数的当前时间curr_date  这里直接传入开始时间
     ***/

    // 日期转换为日期类型
    val dateDF = df.select(
      to_date($"curr_date").as("curr_date"),
      to_date($"end_date").as("end_date"),
      $"p"
    )

    var dateDFDay =dateDF.withColumn("curr_date", date_add($"curr_date",0))

    var currDate= dateDFDay.first().getAs[Date]("curr_date")

    var currDateLastMonth= dateDFDay
      .withColumn("curr_date_last_month", date_format(add_months($"curr_date",-1), "yyyyMM"))
      .withColumn("curr_date_last_month",concat($"p",$"curr_date_last_month"))
      .first().getAs[String]("curr_date_last_month")

    var currDateMonth= dateDFDay
      .withColumn("curr_date_month", date_format(add_months($"curr_date",0), "yyyyMM"))
      .withColumn("curr_date_month",concat($"p",$"curr_date_month"))
      .first().getAs[String]("curr_date_month")

    var currDateNextMonth=  dateDFDay
      .withColumn("curr_date_next_month",date_format(add_months($"curr_date",1), "yyyyMM"))
      .withColumn("curr_date_next_month",concat($"p",$"curr_date_next_month"))
      .first().getAs[String]("curr_date_next_month")

    val endDateDf = dateDF.withColumn("end_date", date_add($"end_date", 0)).first()
    val endDate = endDateDf .getAs[Date]("end_date")

    println(currDate,currDateLastMonth,currDateMonth,currDateNextMonth,endDate)

    println(currDate)


    var j = 0
    while (!currDate.equals(endDate)) {// 到结束日期后结束循环
      println(s"当前值: $j")
        dateDFDay = dateDF.withColumn("curr_date", date_add($"curr_date",j))

        currDate = dateDFDay.first().getAs[Date]("curr_date")

        currDateLastMonth=dateDFDay
          .withColumn("curr_date_last_month", date_format(add_months($"curr_date",-1), "yyyyMM"))
          .withColumn("curr_date_last_month",concat($"p",$"curr_date_last_month"))
          .first().getAs[String]("curr_date_last_month")

        currDateMonth=dateDFDay
          .withColumn("curr_date_month", date_format(add_months($"curr_date",0), "yyyyMM"))
          .withColumn("curr_date_month",concat($"p",$"curr_date_month"))
          .first().getAs[String]("curr_date_month")

        currDateNextMonth=dateDFDay
          .withColumn("curr_date_next_month",date_format(add_months($"curr_date",1), "yyyyMM"))
          .withColumn("curr_date_next_month",concat($"p",$"curr_date_next_month"))
          .first().getAs[String]("curr_date_next_month")

        /***添加所需要循环的传入时间参数的代码***/

        println(currDate.toString,currDateLastMonth,currDateMonth,currDateNextMonth)
      j += 1
    }

/*
   for(i<-0 until 367) {//  367  是为了防止闰年的情况如果循环涉及到超过一整年以上的情况需要加大数值
      if (!currDate.equals(endDate)) {// 到结束日期后结束循环
          dateDFDay = dateDF.withColumn("curr_date", date_add($"curr_date",i))

          currDate = dateDFDay.first().getAs[Date]("curr_date")

          currDateLastMonth=dateDFDay
            .withColumn("curr_date_last_month", date_format(add_months($"curr_date",-1), "yyyyMM"))
            .withColumn("curr_date_last_month",concat($"p",$"curr_date_last_month"))
            .first().getAs[String]("curr_date_last_month")

          currDateMonth=dateDFDay
            .withColumn("curr_date_month", date_format(add_months($"curr_date",0), "yyyyMM"))
            .withColumn("curr_date_month",concat($"p",$"curr_date_month"))
            .first().getAs[String]("curr_date_month")

          currDateNextMonth=dateDFDay
            .withColumn("curr_date_next_month",date_format(add_months($"curr_date",1), "yyyyMM"))
            .withColumn("curr_date_next_month",concat($"p",$"curr_date_next_month"))
            .first().getAs[String]("curr_date_next_month")

        //添加所需要循环的传入时间参数的代码
			   println(currDate.toString,currDateLastMonth,currDateMonth,currDateNextMonth)
      }else{
        println("循环内跳出最终时间:"+currDate)
        spark.stop()
        return  // 跳出整个for循环结束时间循环
      }
    }
*/
    println("循环外最终时间:"+currDate)
    spark.stop()
  }
}
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐