Como calcular e aplicar UDFs vetoriais em Apache Spark: passo a passo
Este tutorial mostra como criar e aplicar UDFs vetoriais em Apache Spark para transformar colunas do tipo array/Vector — útil para pré-processamento de features, normalização e engenharia de dados antes de modelos de ML. Verá como escrever UDFs eficientes, tratar tipos e evitar erros comuns.
Pré-requisitos
- Python 3.8+ e PySpark instalados (local ou cluster).
- Conhecimentos básicos de DataFrame no PySpark.
- Um ficheiro CSV ou DataFrame com uma coluna do tipo array ou Vector.
Passo 1: Porquê usar UDFs vetoriais em Apache Spark
UDFs (User Defined Functions) permitem executar transformações personalizadas. UDFs vetoriais tratam arrays/Vector inteiros de uma só vez (ex.: normalização L2, escala personalizada, transformação de features). São úteis quando não existe função nativa que faça exactamente o que precisa.
Passo 2: Preparar o ambiente e dados de exemplo
Crie um SparkSession e um DataFrame com uma coluna de arrays numéricos. É comum ocorrerem erros de tipo se os elementos não forem float ou se a coluna não for 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)
Passo 3: Escrever uma UDF vetorial simples (normalização L2)
Usaremos pyspark.sql.functions.udf com tipos correctos (ArrayType(DoubleType())). Exemplo: normalizar cada array para norma L2 = 1 — cuidado com arrays nulos ou de norma zero.
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
# garantir floats e evitar divisão por 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)
Passo 4: UDFs vetoriais com NumPy (maior performance local) e cuidado com serialização
NumPy pode acelerar cálculos dentro da UDF, mas aumenta custos de serialização e não paraleliza dentro de cada executor. Use quando operações vetoriais forem pesadas; converta sempre os tipos para 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)
Passo 5: Evitar UDFs quando existir função Spark ML ou SQL nativa
Antes de criar UDFs, verifique se o Spark MLlib não tem transformação nativa (ex.: VectorAssembler, StandardScaler, Normalizer). Funções nativas são muito mais eficientes e evitam erros de serialização. Use UDFs apenas quando necessário.
Passo 6: Tratar tipos Vector do MLlib
Se a sua coluna for Vector (p.ex. DenseVector), converta para lista antes de aplicar a UDF ou escreva uma UDF que aceite Vector. Exemplo de conversão com VectorUDT.
from pyspark.ml.linalg import Vectors, VectorUDT
from pyspark.sql.functions import udf
from pyspark.sql.types import StructType, StructField, IntegerType
# DataFrame com 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 recebe Vector e devolve 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 o resultado
Confirme que a nova coluna existe, tem o tipo esperado e que as normas são 1 (quando aplicável). Utilize show(), printSchema() e algumas verificações numéricas.
df_norm.printSchema()
df_norm.select("features", "features_l2").show(truncate=False)
# verificar norma com uma expressão SQL simples (reconverter para soma de quadrados)
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()
Conclusão
Aprendeu a criar UDFs vetoriais em Apache Spark para normalizar arrays/Vector, quando usar NumPy, como tratar Vector do MLlib e como validar resultados. Próximos passos: explorar StandardScaler/Normalizer do Spark ML para comparar performance ou escrever UDFs em Scala para reduzir overhead. Dica: prefira sempre transformações nativas antes de recorrer a UDFs para melhor performance.