Spark UDF functions in Java
What's inside this article
⌄
- Spark Java UDF example
- How to create Spark UDF Java
- Spark user defined functions Java
- How to change the column in Spark
- How to modify column in Spark
- How to create a user defined function (UDF) and apply it
For example, let’s create a UDF, that takes a String and returns a String.
For Spark < 2.3:
// for Spark < 2.3
UDF1<String, String> myUdf = param -> {
String result = param + "hello_from_udf";
return result;
};
Consider the function declaration: UDF1<String, String> partitionKey.
The first String is the type of the param parameter. The second String is the return type of result.
We also need to register and call the function, see a section below.
For Spark >= 2.3:
// for Spark >= 2.3
UserDefinedFunction myUdf = udf(
(String param) -> {
String result = param + "hello_from_udf";
return result;
}, DataTypes.StringType
);
We also need to register and call the function, see a section below.
UDF Function Registration
After we have declared and described the function, it needs to be registered:
sqlContext.udf().register("myUdfRegisteredName", myUdf, DataTypes.StringType);
UDF Function Call
Then it can be used by calling callUDF:
Dataset<?> df = sqlContext.parquetFile(taskSource);
df.withColumn("calculatedColumn", callUDF("myUdfRegisteredName", col("sourceColumn")))
.write()
.parquet("/data/tmp");
Thus, we will apply myUdf to each value of the sourceColumn, and make the calculatedColumn from that values.