Cómo calcular y aplicar UDFs vectoriales en Apache Spark: paso a paso
Este tutorial muestra cómo crear y aplicar UDFs vectoriales en Apache Spark para transformar columnas del tipo array/Vector — útil para preprocesamiento de features, normalización e ingeniería de datos antes de modelos de ML. Verá cómo escribir UDFs eficientes, tratar tipos y evitar errores comunes.
Pre-requisitos
- Python 3.8+ y PySpark instalados (local o cluster).
- Conocimientos básicos de DataFrame en PySpark.
- Un fichero CSV o DataFrame con una columna del tipo array o Vector.
Paso 1: Por qué usar UDFs vectoriales en Apache Spark
UDFs (User Defined Functions) permiten ejecutar transformaciones personalizadas. UDFs vectoriales tratan arrays/Vector completos de una sola vez (ej.: normalización L2, escala personalizada, transformación de features). Son útiles cuando no existe una función nativa que haga exactamente lo que necesita.
Paso 2: Preparar el entorno y datos de ejemplo
Creé un SparkSession y un DataFrame con una columna de arrays numéricos. Es común que ocurran errores de tipo si los elementos no son float o si la columna no es 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)
Paso 3: Escribir una UDF vectorial simple (normalización L2)
Usaremos pyspark.sql.functions.udf con tipos correctos (ArrayType(DoubleType())). Ejemplo: normalizar cada array para norma L2 = 1 — cuidado con arrays nulos o de norma cero.
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
# garantizar floats y evitar división por cero
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)
Paso 4: UDFs vectoriales con NumPy (mayor rendimiento local) y cuidado con la serialización
NumPy puede acelerar cálculos dentro de la UDF, pero aumenta costes de serialización y no paraleliza dentro de cada executor. Use cuando las operaciones vectoriales sean pesadas; convierta siempre los tipos a listas antes de devolver.
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)
Paso 5: Evitar UDFs cuando exista función Spark ML o SQL nativa
Antes de crear UDFs, verifique si Spark MLlib no tiene una transformación nativa (ej.: VectorAssembler, StandardScaler, Normalizer). Las funciones nativas son mucho más eficientes y evitan errores de serialización. Use UDFs solo cuando sea necesario.
Paso 6: Tratar tipos Vector del MLlib
Si su columna es Vector (p. ej. DenseVector), conviértala a lista antes de aplicar la UDF o escriba una UDF que acepte Vector. Ejemplo de conversión con VectorUDT.
from pyspark.ml.linalg import Vectors, VectorUDT
from pyspark.sql.functions import udf
from pyspark.sql.types import StructType, StructField, IntegerType
# DataFrame con 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 que recibe Vector y devuelve 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)
Verificar el resultado
Confirme que la nueva columna existe, tiene el tipo esperado y que las normas son 1 (cuando proceda). Utilice show(), printSchema() y algunas comprobaciones numéricas.
df_norm.printSchema()
df_norm.select("features", "features_l2").show(truncate=False)
# verificar norma con una expresión SQL simple (reconvertir a suma de cuadrados)
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()
Conclusión
Ha aprendido a crear UDFs vectoriales en Apache Spark para normalizar arrays/Vector, cuándo usar NumPy, cómo tratar Vector del MLlib y cómo validar resultados. Siguientes pasos: explorar StandardScaler/Normalizer de Spark ML para comparar rendimiento o escribir UDFs en Scala para reducir overhead. Consejo: prefiera siempre transformaciones nativas antes de recurrir a UDFs para mejor rendimiento.