Failed to execute UDF

Viewed 102

there is the table in which the field "A" contains sql query. It is necessary to add an additional field "B" that would contain the time spent on the execution of the query from the "A" field. I wrote a UDF and everything works well, but when caching the resulting table or trying to write the final dataframe to a physical table, I got the error:

"Failed to execute user defined function ($anonfun$1: (string) => string)"

. What could be the problem? Example:

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
})

Source table "source":

+--------------------+
|          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  |
+--------------------+------+

But:

//ERROR
result.write.mode("overwrite").saveAsTable("dbName.result")

//ERROR
val result_cache = result.persist
result_cache.show
1 Answers

The issue here is that UDF works on a executors where spark session isn't available. So I guess you get NullPointer exception on a "val count = spark.sql..." line.

You should do it on a driver using not a UDF but just function1. Also using collect() I suppose that the main table isn't big and will fit into a driver memory:

Example:

import java.util.Date
import java.time.LocalDateTime

val set_time = (query: String) => {
val start = new Timestamp(new Date().getTime)
val count = spark.sql(s"${query}").count
val time_query = (new Timestamp(new Date().getTime)).getTime() - start.getTime()
time_query.toString
}

val result = spark.sql("select 'select 1' as A union all select 'select 2' as A")
val s = result.collect().map(x =>(x(0).asInstanceOf[String],set_time(x(0).asInstanceOf[String]))).toList.toDF("A","B")

s.show 
s.cache().show 
+--------+---+
|       A|  B|
+--------+---+
|select 1|171|
|select 2|135|
+--------+---+

PS: also val start = new Timestamp(new Date().getDate) in your example should be val start = new Timestamp(new Date().getTime)

Related