在Spark DataFrame中使用Scala获取下一个工作日的实现方案
Custom UDF for Next Workday in Spark 2.0.2
Since Spark 2.0.2 doesn’t include a ready-made function for this exact date logic, we can build a User-Defined Function (UDF) to meet your requirements: return the next calendar day if the input is a weekday (Mon-Fri), or return the following Monday if the input is Saturday or Sunday.
Implementation Options
Scala Version
First, import the necessary dependencies, then define and use the UDF:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types.DateType import java.sql.Date import java.time.{LocalDate, DayOfWeek} // Define the UDF val nextWorkdayUdf = udf((inputDate: Date) => { val localDate = inputDate.toLocalDate val nextDate = localDate.getDayOfWeek match { // Saturday: add 2 days to reach Monday case DayOfWeek.SATURDAY => localDate.plusDays(2) // Sunday: add 1 day to reach Monday case DayOfWeek.SUNDAY => localDate.plusDays(1) // Weekdays: just add 1 day case _ => localDate.plusDays(1) } Date.valueOf(nextDate) }) // Use the UDF on your DataFrame // Assume your DataFrame is named `df` with a DateType column `input_date` val resultDF = df.withColumn("next_workday", nextWorkdayUdf(col("input_date")))
Java Version
If you’re working with the Java API, here’s the equivalent implementation:
import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.api.java.UDF1; import org.apache.spark.sql.functions; import org.apache.spark.sql.types.DataTypes; import java.sql.Date; import java.time.LocalDate; import java.time.DayOfWeek; public class NextWorkdayUdfExample { public static void main(String[] args) { SparkSession spark = SparkSession.builder() .appName("NextWorkdayUdf") .master("local[*]") .getOrCreate(); // Define the UDF UDF1<Date, Date> nextWorkdayUdf = (inputDate) -> { LocalDate localDate = inputDate.toLocalDate(); LocalDate nextDate; switch (localDate.getDayOfWeek()) { case SATURDAY: nextDate = localDate.plusDays(2); break; case SUNDAY: nextDate = localDate.plusDays(1); break; default: nextDate = localDate.plusDays(1); break; } return Date.valueOf(nextDate); }; // Register the UDF with Spark spark.udf().register("next_workday", nextWorkdayUdf, DataTypes.DateType); // Use the UDF in your DataFrame // Assume your DataFrame is named `df` with a DateType column `input_date` spark.sql("SELECT input_date, next_workday(input_date) AS next_workday FROM df").show(); // Or using DataFrame API // df.withColumn("next_workday", functions.callUDF("next_workday", functions.col("input_date"))).show(); } }
Legacy Java Version (Pre-Java 8)
If your environment doesn’t support Java 8’s java.time API, you can use java.util.Calendar instead:
// Scala example with Calendar val nextWorkdayUdf = udf((inputDate: Date) => { val cal = java.util.Calendar.getInstance() cal.setTime(inputDate) val dayOfWeek = cal.get(java.util.Calendar.DAY_OF_WEEK) val addDays = dayOfWeek match { case java.util.Calendar.SATURDAY => 2 case java.util.Calendar.SUNDAY => 1 case _ => 1 } cal.add(java.util.Calendar.DATE, addDays) new Date(cal.getTimeInMillis()) })
Verification with Your Examples
- Input
2017-01-01(Monday, weekday): UDF adds 1 day → returns2017-01-02(Tuesday) ✔️ - Input
2017-04-03(Saturday): UDF adds 2 days → returns2017-04-05(Monday) ✔️
内容的提问来源于stack exchange,提问作者user2441441
相关产品推荐
相关产品推荐

