Saltar al contenido

DataFrames

Menú

Puedes leer sin crear un espacio. Créalo solo cuando quieras guardar.

Guardar progreso

DataFrames, transformaciones y datos complejos

Domina las operaciones que aparecen en ETL real: joins, arrays, structs, ventanas y deduplicación.

Lectura pública
Detalles
Reto observable

Dado un dataset semiestructurado, publica una tabla normalizada con esquema explícito y demuestra unicidad de claves, cardinalidad y tratamiento de registros inválidos.

Al terminar podrás
  • Transformar columnas y filas con funciones nativas
  • Manipular arrays, maps y structs
  • Combinar y deduplicar datasets de forma determinista
Prerrequisitos
m03
Última revisión
25 ago 2026
Nivel
Associate
Ruta relacionada
core
Dominios blueprint
Data Transformation and Modeling · PySpark DataFrames
Estado
Revisión editorial interna
Fuentes principales
PySpark DataFrames · DataFrame transformations
Reportar un error
01
Modelo mental

Transformaciones y acciones

Las transformaciones construyen un plan lógico; las acciones lo ejecutan y materializan un resultado.

DataFrames son inmutables: select, filter o withColumn devuelven un nuevo plan sin leer todos los datos inmediatamente. Catalyst puede reorganizar y simplificar ese plan antes de ejecutarlo.

count, collect, write y display son acciones. Repetir acciones sobre el mismo linaje puede recalcularlo; cache solo compensa si hay reutilización medida y memoria suficiente.

PySparkObservar lazy evaluation
paid = orders.filter("status = 'paid'").select("order_id", "amount")
paid.explain("formatted")
rows = paid.count()

explain inspecciona el plan; count inicia la ejecución.

¿filter ejecuta inmediatamente una lectura?

Profundiza

La API DataFrame es perezosa: una transformación describe una nueva relación y una acción exige un resultado. Spark no procesa fila por fila al escribir select, filter o join; acumula un plan lógico que Catalyst puede reorganizar. count, collect, write o una visualización desencadenan un job que materializa parte del linaje. Esta separación permite pushdown, poda de columnas y elección de joins, pero también sorprende cuando una acción repetida recalcula todo. Cache solo tiene sentido si el mismo resultado costoso se reutiliza y cabe con seguridad. Para razonar sobre rendimiento, identifica dónde se define el plan, dónde nace una acción y qué fronteras de shuffle dividen stages.

Transformación

Operación perezosa que produce un nuevo DataFrame y amplía el plan lógico.

Permite componer trabajo antes de que Spark elija cómo ejecutarlo.
Acción

Operación que solicita un resultado y desencadena la ejecución del plan necesario.

Marca el punto donde aparecen coste, jobs y errores ligados a datos.
Linaje

Cadena de dependencias de transformaciones necesaria para recomputar un DataFrame.

Explica recalculo, recuperación y cuándo una caché puede aportar valor.
Resumen

Puntos clave

  • Transformaciones son lazy
  • Acciones disparan jobs
  • Cache es una decisión medida

Evita

  • Usar collect sobre un dataset grande
  • Cachear cada DataFrame por costumbre
02
Implementación

Select, filter y withColumn

La limpieza fiable hace explícitos nulos, tipos, nombres y reglas antes de publicar Silver.

Convierte tipos con cast o try_cast según la política de error, normaliza texto con funciones nativas y decide si un nulo se rechaza, imputa o conserva. Una conversión silenciosa a null debe medirse.

select con expresiones explícitas produce contratos más revisables que arrastrar todas las columnas. Añade columnas técnicas como ingestion_ts o source_file cuando aporten trazabilidad.

SQLLimpieza defensiva
SELECT order_id,
       try_cast(amount AS DECIMAL(18,2)) AS amount,
       lower(trim(email)) AS email
FROM main.bronze.orders
WHERE order_id IS NOT NULL;

Cuenta los amount convertidos a null antes de aceptar el lote.

¿Qué ventaja ofrece try_cast?

Profundiza

Limpiar datos significa convertir ambigüedad de origen en un contrato explícito, no encadenar dropna y cast hasta que el job termine. Primero se define el grain de la tabla, las columnas canónicas, tipos, nulos permitidos, zonas horarias y reglas de dominio. Después se distinguen tres resultados: válido, corregible y rechazado. try_cast convierte errores de representación en null para poder medirlos; cast estricto puede ser preferible cuando el contrato exige fallo. Silver debe conservar claves y evidencia de origen suficientes para reconciliar. Una regla de calidad sin denominador, umbral y acción es solo una expresión, no un control operativo.

Grain

Significado exacto de una fila y conjunto mínimo de dimensiones que identifica un hecho.

Evita duplicados conceptuales y agregaciones incorrectas en capas posteriores.
try_cast

Conversión que devuelve null cuando un valor no puede representarse en el tipo solicitado.

Permite cuantificar errores sin abortar todo el lote, siempre que se controle el null resultante.
Cuarentena

Destino gobernado para registros inválidos junto con su causa y contexto de ingestión.

Conserva evidencia y permite reparación sin contaminar la tabla confiable.
Resumen

Puntos clave

  • Cada nulo requiere una política
  • El esquema objetivo debe ser explícito
  • Las funciones nativas conservan optimización

Evita

  • Rellenar todos los nulos con cero
  • Usar SELECT * en un contrato Silver
03
Operación

Joins, unions y claves compuestas

El tipo de join expresa qué filas deben sobrevivir; la cardinalidad y las claves determinan corrección y coste.

Inner conserva coincidencias; left conserva todas las filas de la izquierda; full conserva ambos lados. Antes del join comprueba unicidad: una dimensión duplicada puede multiplicar hechos sin producir error.

union combina por posición y unionByName por nombre; ninguna elimina duplicados. Broadcast puede evitar shuffle para un lado pequeño, pero solo tras confirmar tamaño.

PySparkJoin con claves compuestas
enriched = orders.join(
    customers.select("tenant_id", "customer_id", "segment"),
    ["tenant_id", "customer_id"],
    "left",
)

Compara row count y claves sin correspondencia después del join.

¿Qué puede revelar un aumento inesperado de filas?

Profundiza

Un join combina conjuntos según predicados, pero su corrección depende del grain y la cardinalidad antes que del tipo sintáctico. Inner conserva coincidencias; left preserva todas las filas izquierdas; semi responde existencia sin añadir columnas; anti conserva ausencias. Si una clave es única en un lado y repetida en otro, el resultado puede multiplicar filas legítimamente. Si se esperaba uno a uno, esa multiplicación es un defecto de datos o de predicado. Antes de optimizar broadcast o particiones, declara qué fila representa cada entrada, normaliza claves y estima conteos. Un join rápido que duplica ingresos es peor que uno lento: semántica y reconciliación son criterios de aceptación.

Cardinalidad

Relación de multiplicidad entre claves de dos datasets, como uno-a-uno o uno-a-muchos.

Predice el número de filas y revela duplicaciones accidentales.
Left semi join

Join que conserva filas izquierdas con al menos una coincidencia sin añadir columnas derechas.

Es la forma precisa y eficiente de filtrar por existencia.
Null-safe equality

Comparación que considera dos null equivalentes mediante una semántica explícita.

Evita asumir que la igualdad SQL ordinaria empareja valores desconocidos.
Resumen

Puntos clave

  • Elige join por semántica
  • Valida cardinalidad antes y después
  • Union no deduplica

Evita

  • Unir solo por customer_id en un sistema multitenant
  • Confundir union con eliminación de duplicados
04
Diagnóstico

explode, arrays, maps y structs

Arrays, structs, maps y VARIANT cubren contratos distintos de datos complejos sin perder el contexto de la fila padre.

explode crea una fila por elemento; explode_outer conserva la fila cuando la colección es null o vacía. Los campos de un struct se seleccionan con notación de punto y transform procesa arrays sin expandirlos. Para un esquema conocido que necesita estadísticas y layout, un struct tipado sigue siendo la opción más explícita.

VARIANT permite almacenar JSON semiestructurado con codificación optimizada en Databricks Runtime 15.3 o superior cuando los campos cambian con frecuencia. No sustituye el modelado: las columnas VARIANT no son claves de particionado o clustering ni se agrupan u ordenan directamente; proyecta a columnas tipadas los campos usados por contratos y consultas críticas y verifica compatibilidad de runtime/protocolo.

PySparkNormalizar líneas de pedido
lines = (raw
  .select("order_id", F.explode_outer("items").alias("item"))
  .select("order_id", F.col("item.sku").alias("sku"), F.col("item.qty").alias("qty")))

Decide si una fila con sku null debe ir a cuarentena.

¿Cuándo elegir explode_outer?

Profundiza

Los tipos complejos conservan estructura: un struct agrupa campos con esquema, un array mantiene una secuencia y un map asocia claves con valores. No es necesario convertir JSON a cadenas ni explotar todo inmediatamente. Las funciones de orden superior transform, filter, exists y aggregate operan dentro de un array preservando la fila padre; la notación de campo navega structs; element_at consulta colecciones. explode cambia el grain al crear filas y, por tanto, exige conservar claves del padre. explode_outer mantiene una representación cuando la colección es null o vacía. Elegir entre transformación anidada y normalización depende del consumidor y de la semántica, no de una limitación de Spark.

Struct

Valor compuesto con campos nombrados y tipos definidos dentro de una columna.

Permite conservar jerarquía y seleccionar atributos sin perder el contrato.
Función de orden superior

Expresión que aplica una operación a elementos de una colección sin convertirlos en filas independientes.

Preserva grain y suele simplificar transformaciones de arrays.
explode_outer

Generador que expande elementos y conserva el padre con null cuando la colección no aporta elementos.

Evita perder entidades padre cuando la ausencia de detalles tiene significado.
Resumen

Puntos clave

  • explode cambia cardinalidad
  • Struct favorece contrato conocido
  • VARIANT conserva flexibilidad; proyecta campos críticos

Evita

  • Perder order_id al explotar items
  • Guardar todo como VARIANT y esperar estadísticas, clustering o GROUP BY directos
05
Decisión de diseño

Ventanas, agregaciones y deduplicación

Ventanas calculan métricas por grupo sin colapsar filas y permiten deduplicación determinista.

groupBy reduce cada grupo; una window mantiene cada registro y añade ranking, acumulados o valores previos. El orden debe resolver empates para que el resultado sea repetible.

Para conservar la versión más reciente, combina row_number con una ordenación por event_ts y un segundo campo estable. dropDuplicates sin criterio temporal no expresa qué versión conservar.

PySparkÚltimo evento por clave
w = Window.partitionBy("order_id").orderBy(F.desc("event_ts"), F.desc("ingest_id"))
latest = events.withColumn("rn", F.row_number().over(w)).filter("rn = 1").drop("rn")

ingest_id rompe empates de event_ts y hace la salida determinista.

¿Por qué dropDuplicates no basta para elegir el evento más nuevo?

Profundiza

Una window calcula sobre un conjunto relacionado con cada fila sin colapsarlo como groupBy. PARTITION BY define el grupo lógico, ORDER BY establece secuencia y el frame delimita qué filas contribuyen. row_number asigna una prioridad total solo si el orden contiene un desempate estable; rank y dense_rank expresan empates con semánticas diferentes. Para deduplicar, primero se define la identidad del evento y qué versión debe ganar; después se ordena por tiempo de negocio, secuencia y un identificador determinista. dropDuplicates expresa igualdad, no preferencia, y en batch no garantiza conservar el evento más nuevo. El resultado se valida por unicidad y reconciliación de versiones descartadas.

Window specification

Definición de partición, orden y frame utilizada para calcular una función por cada fila.

Determina tanto la semántica como el movimiento de datos de una window.
row_number

Función que asigna una secuencia única dentro de cada partición según el orden declarado.

Permite seleccionar un único ganador cuando el orden es totalmente determinista.
Desempate estable

Columna adicional única o de orden consistente usada cuando el criterio principal empata.

Evita que reintentos elijan versiones diferentes con los mismos timestamps.
Resumen

Puntos clave

  • Window no colapsa filas
  • El orden debe ser total
  • Deduplicar exige regla de supervivencia

Evita

  • Omitir criterio de desempate
  • Usar groupBy cuando se necesitan columnas de detalle

Fuente revisada · vista externa

Azure Databricks Hands-on

Azure Databricks Hands-on · commit a91650b

HandsOn.dbc

Archivo importable

Este notebook se abre desde su fuente revisada

El repositorio no permite republicar su contenido dentro de Lakehouse Lab. Conservamos la misma experiencia lateral, la ruta exacta y el commit auditado, y dejamos la lectura en GitHub para respetar la autoría.

Autor
Tsuyoshi Matsuzaki
Licencia
No verificada
Formato
dbc
Abrir / descargar .dbc

Módulo 04

Contenido del módulo