Spark Scala中TimestampType重置分秒的等价方法及时间差计算
Got it, let's break this down for you—first, a quick correction: in your UDF definition, the parameter types should be java.sql.Timestamp (the runtime type Spark uses for TimestampType columns) instead of org.apache.spark.sql.types.TimestampType (that's just schema metadata, not an instantiable class).
Option 1: Use Java 8+ DateTime API in UDF
Scala doesn't have a direct replace method like Python for timestamps, but we can leverage Java's LocalDateTime API to achieve the same truncation effect. Here's how to adjust your UDF:
import java.sql.Timestamp import java.time.LocalDateTime import org.apache.spark.sql.functions.udf val split_hour_range_udf = udf { (startDateTime: Timestamp, endDateTime: Timestamp) => // Convert Timestamp to LocalDateTime for easy manipulation val startLocal = startDateTime.toLocalDateTime // Reset minutes, seconds, and nanoseconds to 0 to get the hour-aligned time val truncatedStart = startLocal.withMinute(0).withSecond(0).withNano(0) val adjustedStart = Timestamp.valueOf(truncatedStart) // Repeat the truncation for the end time val endLocal = endDateTime.toLocalDateTime val truncatedEnd = endLocal.withMinute(0).withSecond(0).withNano(0) val adjustedEnd = Timestamp.valueOf(truncatedEnd) // Calculate time difference (example: return total hours as a Long) val timeDiffMillis = adjustedEnd.getTime - adjustedStart.getTime timeDiffMillis / (1000L * 60 * 60) // Convert milliseconds to hours }
Option 2: Use Spark Built-in Function (Better Performance)
If you don't strictly need a UDF, Spark's date_trunc function is optimized and avoids the overhead of UDFs. This is the preferred approach for most cases, as Spark can optimize built-in functions much better:
import org.apache.spark.sql.functions.{date_trunc, col} import org.apache.spark.sql.types.LongType val resultDF = statusWithOutDuplication // Fix your timestamp format first—note lowercase 'yyyy-MM-dd HH:mm:ss' (Spark follows Java's SimpleDateFormat rules) .withColumn("requestTime", unix_timestamp(col("requestTime"), "yyyy-MM-dd HH:mm:ss").cast(TimestampType)) .withColumn("responseTime", unix_timestamp(col("responseTime"), "yyyy-MM-dd HH:mm:ss").cast(TimestampType)) // Truncate timestamps to the nearest hour (automatically resets min/sec to 0) .withColumn("truncatedRequest", date_trunc("hour", col("requestTime"))) .withColumn("truncatedResponse", date_trunc("hour", col("responseTime"))) // Calculate the hour difference .withColumn("hourDifference", (col("truncatedResponse").cast(LongType) - col("truncatedRequest").cast(LongType)) / (1000L * 60 * 60))
One quick note on your original timestamp format: YYYY-MM-DD HH:MM:SS should be yyyy-MM-dd HH:mm:ss—uppercase MM is for months, lowercase mm for minutes; uppercase DD is day-of-year, lowercase dd is day-of-month.
内容的提问来源于stack exchange,提问作者syv

