You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Scala中TimestampType重置分秒的等价方法及时间差计算

Solution for Truncating Timestamps in Scala Spark UDF

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 03:50:49