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.