【发布时间】:2020-06-13 13:51:26
【问题描述】:
我有以下 UDF 用于将存储为字符串的时间转换为时间戳。
val hmsToTimeStampUdf = udf((dt: String) => {
if (dt == null) null else {
val formatter = DateTimeFormat.forPattern("HH:mm:ss")
try {
new Timestamp(formatter.parseDateTime(dt).getMillis)
} catch {
case t: Throwable => throw new RuntimeException("hmsToTimeStampUdf,dt="+dt, t)
}
}
})
这个UDF用于将String值转换成Timestamp:
outputDf.withColumn(schemaColumn.name, ymdToTimeStampUdf(col(schemaColumn.name))
但某些 CSV 文件的此列的值无效,导致 RuntimeException。我想找出哪些行有这些损坏的记录。是否可以访问 UDF 中的行信息?
【问题讨论】:
-
嗯,不,但你可以有2个参数的UDF,时间列第二是ID列,然后打印它或其他东西。这个策略对你有用吗?
标签: apache-spark user-defined-functions