Spark DataFrame中Scala UDF无法运行问题求助
问题原因
你的报错源于Spark UDF无法识别Scala函数的默认参数:原函数convertNanoEpochToDateTime定义了4个参数(其中3个带默认值),但用udf(convertNanoEpochToDateTime _)包装后,生成的UDF会要求传入全部4个参数,而你调用时只传了1个时间列,导致参数数量不匹配。
修复方案
方案1:封装单参数函数(推荐)
把默认参数硬编码到新函数中,再包装成UDF,调用时只需传入时间列:
import java.text.SimpleDateFormat import java.util.{Date, TimeZone} // 重新定义仅接收Long类型的函数,内置默认配置 def convertNanoEpochToCET(d: Long): String = { val format = "dd/MM/yyyy HH:mm:ss.SSS" val timeZone = "CET" val msPrecision = 9 val sdf = new SimpleDateFormat(format) sdf.setTimeZone(TimeZone.getTimeZone(timeZone)) // 纳秒转毫秒:直接除以1e6,避免浮点运算误差 val date = new Date(d / 1000000L) val stringTime = sdf.format(date) if (format.contains(".S")) { // 提取纳秒部分的后9位 val milliSecondsStr = d.toString.takeRight(9) stringTime.substring(0, stringTime.lastIndexOf(".") + 1) + milliSecondsStr.substring(0, msPrecision) } else { stringTime } } // 包装为UDF val epochToCET = udf(convertNanoEpochToCET _)
调用方式不变:
val df2 = df1.select($"messageID", $"messageIndex", epochToCET($"messageTimestamp").as("messageTimestamp"))
方案2:显式传递默认参数
如果需要保留函数的灵活性,调用UDF时用lit()包装常量,显式传入所有默认参数:
// 原函数和UDF定义不变 val epochToDateTime = udf(convertNanoEpochToDateTime _) // 调用时传入全部参数 val df2 = df1.select( $"messageID", $"messageIndex", epochToDateTime($"messageTimestamp", lit("dd/MM/yyyy HH:mm:ss.SSS"), lit("CET"), lit(9)).as("messageTimestamp") )
进阶优化:线程安全的日期API
SimpleDateFormat不是线程安全的,在Spark分布式环境下可能出现异常,推荐改用Java 8的DateTimeFormatter:
import java.time.{Instant, ZoneId} import java.time.format.DateTimeFormatter def convertNanoEpochToCET(d: Long): String = { val format = "dd/MM/yyyy HH:mm:ss.SSS" val timeZone = ZoneId.of("CET") val msPrecision = 9 // 直接用纳秒构建Instant,精度更准确 val instant = Instant.ofEpochSecond(d / 1000000000L, d % 1000000000L) val zonedDateTime = instant.atZone(timeZone) val formatter = DateTimeFormatter.ofPattern(format) val baseTime = formatter.format(zonedDateTime) if (format.contains(".S")) { // 处理纳秒截断,补零保证长度 val nanoPart = (d % 1000000000L).toString.reverse.padTo(9, '0').reverse val truncatedNano = nanoPart.substring(0, msPrecision) baseTime.substring(0, baseTime.lastIndexOf(".") + 1) + truncatedNano } else { baseTime } } val epochToCET = udf(convertNanoEpochToCET _)
内容的提问来源于stack exchange,提问作者Vijay B
相关产品推荐
相关产品推荐

