How to calculate and apply vector UDFs in Apache Spark: step by step
This tutorial shows how to create and apply vector UDFs in Apache Spark to transform columns of type array/Vector — useful for feature preprocessing, normalization and data engineering before ML models. You will see how to write efficient UDFs, handle types and avoid common errors.
Prerequisites
- Python 3.8+ and PySpark installed (local or cluster).
- Basic knowledge of DataFrame in PySpark.
- A CSV file or DataFrame with a column of type array or Vector.
Step 1: Why use vector UDFs in Apache Spark
UDFs (User Defined Functions) allow executing custom transformations. Vector UDFs handle entire arrays/Vectors at once (e.g., L2 normalization, custom scaling, feature transformation). They are useful when there is no native function that does exactly what you need.
Step 2: Prepare the environment and sample data
Create a SparkSession and a DataFrame with a column of numeric arrays. Type errors commonly occur if elements are not float or if the column is not ArrayType.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col
spark = SparkSession.builder.appName("udf-vetorial-exemplo").getOrCreate()
data = [ (1, [1.0, 2.0, 3.0]), (2, [0.0, 0.0, 0.0]), (3, [4.0, 5.0, 6.0]) ]
df = spark.createDataFrame(data, ["id", "features"])
df.printSchema()
df.show(truncate=False)
Step 3: Write a simple vector UDF (L2 normalization)
We will use pyspark.sql.functions.udf with correct types (ArrayType(DoubleType())). Example: normalize each array to L2 norm = 1 — be careful with null arrays or zero norm.
from pyspark.sql.functions import udf
from pyspark.sql.types import ArrayType, DoubleType
import math
def l2_normalize(arr):
if arr is None:
return None
# ensure floats and avoid division by zero
try:
vals = [float(x) for x in arr]
except Exception:
return None
norm = math.sqrt(sum(x*x for x in vals))
if norm == 0.0:
return [0.0 for _ in vals]
return [x / norm for x in vals]
l2_udf = udf(l2_normalize, ArrayType(DoubleType()))
df_norm = df.withColumn("features_l2", l2_udf(col("features")))
df_norm.show(truncate=False)
Step 4: Vector UDFs with NumPy (higher local performance) and serialization caution
NumPy can speed up calculations inside the UDF, but increases serialization costs and does not parallelize within each executor. Use when vector operations are heavy; always convert types to lists before returning.
import numpy as np
def l2_normalize_numpy(arr):
if arr is None:
return None
a = np.array(arr, dtype=float)
norm = np.linalg.norm(a)
if norm == 0.0:
return [0.0]*len(a)
return (a / norm).tolist()
l2_udf_np = udf(l2_normalize_numpy, ArrayType(DoubleType()))
df.withColumn("features_l2_np", l2_udf_np(col("features"))).show(truncate=False)
Step 5: Avoid UDFs when there is a Spark ML or native SQL function
Before creating UDFs, check whether Spark MLlib has a native transformation (e.g., VectorAssembler, StandardScaler, Normalizer). Native functions are much more efficient and avoid serialization errors. Use UDFs only when necessary.
Step 6: Handling MLlib Vector types
If your column is Vector (e.g. DenseVector), convert to a list before applying the UDF or write a UDF that accepts Vector. Example conversion with VectorUDT.
from pyspark.ml.linalg import Vectors, VectorUDT
from pyspark.sql.functions import udf
from pyspark.sql.types import StructType, StructField, IntegerType
# DataFrame with Vector
data2 = [ (1, Vectors.dense([1.0,2.0,3.0])), (2, Vectors.dense([0.0,0.0,0.0])) ]
df_vec = spark.createDataFrame(data2, ["id","features_vec"])
# UDF that receives Vector and returns ArrayType(DoubleType())
def vec_to_norm(v):
if v is None:
return None
arr = v.toArray().tolist()
return l2_normalize(arr)
vec_udf = udf(vec_to_norm, ArrayType(DoubleType()))
df_vec.withColumn("features_l2", vec_udf(col("features_vec"))).show(truncate=False)
Verify the result
Confirm that the new column exists, has the expected type and that the norms are 1 (when applicable). Use show(), printSchema() and some numerical checks.
df_norm.printSchema()
df_norm.select("features", "features_l2").show(truncate=False)
# check norm with a simple SQL expression (reconvert to sum of squares)
from pyspark.sql.functions import expr
check = df_norm.withColumn("norm_sq", expr("aggregate(features_l2, 0D, (acc, x) -> acc + x*x)"))
check.select("id", "norm_sq").show()
Conclusion
You learned how to create vector UDFs in Apache Spark to normalize arrays/Vector, when to use NumPy, how to handle MLlib Vector and how to validate results. Next steps: explore Spark ML StandardScaler/Normalizer to compare performance or write UDFs in Scala to reduce overhead. Tip: always prefer native transformations before resorting to UDFs for better performance.