【发布时间】:2020-12-24 15:49:51
【问题描述】:
存在字段“A”包含sql查询的表。有必要添加一个附加字段“B”,该字段将包含从“A”字段执行查询所花费的时间。我写了一个 UDF,一切正常,但是在缓存结果表或尝试将最终数据帧写入物理表时,出现错误:
"执行用户定义函数失败 ($anonfun$1: (string) => 字符串)"
。可能是什么问题呢? 示例:
val set_time = udf((query: String) => {
val start = new Timestamp(new Date().getDate)
val count = spark.sql(s"${query}").count
val time_query = (new Timestamp(new Date().getTime)).getTime() - start.getTime()
time_query.toString
})
源表“源”:
+--------------------+
| A |
+--------------------+
|"Select * From ..." |
|"Select * From ..." |
|"Select * From ..." |
|"Select * From ..." |
|"Select * From ..." |
+--------------------+
val result = spark.sql("from source").
withColumn("B", set_time(col("A")))
result.show
+--------------------+------+
| A | B |
+--------------------+------+
|"Select * From ..." | 356 |
|"Select * From ..." | 642 |
|"Select * From ..." | 2745 |
|"Select * From ..." | 1324 |
|"Select * From ..." | 635 |
+--------------------+------+
但是:
//ERROR
result.write.mode("overwrite").saveAsTable("dbName.result")
//ERROR
val result_cache = result.persist
result_cache.show
【问题讨论】:
标签: scala apache-spark user-defined-functions