DP-700 · MÓDULO 3 🔥

Usar Apache Spark en Microsoft Fabric

El motor de cómputo del data engineering: arquitectura, pools, nodos, runtimes, environments, optimizaciones, notebooks vs Spark Job Definitions, DataFrames, particionado, Spark SQL y visualización. 🌸

Intermedio Avanzado
🎯

1. Objetivos y encaje en el DP-700

¿Qué te enseña este módulo?

Este es el módulo donde aprendemos a usar el motor de cómputo del data engineering en Fabric. Los objetivos oficiales:

  1. Configurar Spark en un workspace de Fabric.
  2. Identificar escenarios adecuados para notebooks Spark y Spark jobs.
  3. Usar Spark dataframes para analizar y transformar datos.
  4. Usar Spark SQL para consultar datos en tablas y vistas.
  5. Visualizar datos en un notebook Spark.

Peso en el examen DP-700

Este módulo es fundamental para el Dominio 2 (Ingest and transform) y aparece transversalmente en el resto:

DominioCómo aparece Spark
Implement and manage (30-35%)Configuración de pools, runtimes y environments a nivel workspace
Ingest and transform (30-35%)PySpark y Spark SQL para transformar en notebooks, patrones de escritura Delta
Monitor and optimize (30-35%)Optimizaciones de Spark (autoscale, dynamic allocation, native execution engine, autotune)

🎯 Qué es NUEVO respecto a la DP-600

En DP-600 se toca Spark de pasada (crear un notebook, hacer un spark.read.load()). En DP-700 hay que dominar la configuración, los conceptos de rendimiento y los matices de PySpark. Cosas que salen que antes no:

  • 🆕 Native Execution Engine (vectorized processing)
  • 🆕 High Concurrency Mode (compartir sesiones)
  • 🆕 Environments (custom Spark configs por workload)
  • 🆕 Runtime 1.3 vs 2.0 (versiones de Spark/Delta)
  • 🆕 Elección Python kernel vs Spark kernel
  • 🆕 NotebookUtils (antes MSSparkUtils)
  • 🆕 Autotune para optimización automática
🔥

2. ¿Qué es Apache Spark?

Definición

Apache Spark es un framework de procesamiento distribuido de datos open-source. Su gran idea: cuando tienes muchísimos datos que no caben ni se procesan bien en una sola máquina, Spark distribuye el trabajo entre varias máquinas (nodos) y coordina la ejecución.

En Fabric, ese conjunto de máquinas se llama un Spark pool.

La filosofía "divide y vencerás"

Piensa en Spark como una fábrica de galletas kawaii 🍪:

  • Si tienes que hacer 10 galletas → una persona sola, en su cocina, las hace.
  • Si tienes que hacer 1 millón de galletas → necesitas repartir el trabajo entre muchos operarios en muchas cocinas trabajando en paralelo.

Spark es "el gerente" que decide qué operario hace cada parte, recoge los resultados y te devuelve las galletas terminadas. 🌸

Lenguajes soportados

Spark ejecuta código escrito en varios lenguajes:

  • PySpark (Python) → el más popular
  • Spark SQL → sintaxis SQL clásica
  • Scala → el "nativo" de Spark
  • SparkR → R
  • Java → el más verboso, poco común

En la práctica, PySpark y Spark SQL cubren el 95% de los casos.

🎯 Notación clave

Cuando escribes código Spark, casi siempre acabas invocando el objeto spark (una SparkSession). Es el punto de entrada a todo:

spark  # <-- Objeto SparkSession, ya está creado automáticamente en cada notebook

Desde spark puedes leer, escribir, ejecutar SQL, gestionar el catálogo, etc.

🏗️

3. Arquitectura: driver, workers y executors

Este apartado es más "conceptual" pero cae en preguntas de troubleshooting.

Los tres tipos de "procesos" en Spark

A) Driver (uno solo)

  • Vive en el head node.
  • Es el "cerebro": recibe tu código, lo compila, planifica cómo distribuirlo, y coordina.
  • Contiene la SparkSession (el objeto spark).
  • Si el driver se cae → la sesión entera se cae.

B) Workers (varios)

  • Son las máquinas que hacen el trabajo real.
  • Cada worker es un nodo del pool.

C) Executors (uno o varios por worker)

  • Son los procesos que ejecutan tareas dentro de cada worker.
  • Pueden ejecutar múltiples tareas en paralelo (dependen de cores).

🎯 Detalle crítico en Fabric

En Fabric, la ratio nodos a executors es siempre 1:1. Cuando configuras un pool con N nodos, uno se reserva para el driver y los demás son executors.

Excepción: en configuración single-node, el driver y el executor comparten recursos (divididos por 2 aprox.).

Diagrama mental

Head Node                    Worker Node 1        Worker Node 2
┌──────────────┐             ┌─────────────┐      ┌─────────────┐
│   DRIVER     │◄────────────┤  Executor 1 │      │  Executor 2 │
│              │             │             │      │             │
│  SparkSession│─────────────┤  Executor 1 │      │  Executor 2 │
│              │             │             │      │             │
│  Cluster mgr │─────────────┤   Tasks     │      │   Tasks     │
└──────────────┘             └─────────────┘      └─────────────┘
🌊

4. Spark pools en Fabric

Un Spark pool es el conjunto de nodos de cómputo que ejecuta el trabajo Spark. Fabric ofrece dos tipos:

4.1 · Starter pool

Un starter pool se crea automáticamente en cada workspace de Fabric. Es un pool pre-hidratado ("live pool"), lo que significa que los nodos ya están calentitos esperando.

Ventajas:

  • ⚡ Arranque muy rápido (~segundos) porque los nodos están live.
  • 🚀 Ideal para desarrollo, exploración, casos ad-hoc.
  • 💰 Compartido entre workloads.

Limitaciones:

  • Configuración limitada (medium size por defecto).
  • No permite ciertas configuraciones avanzadas como managed private endpoint o private link.

4.2 · Custom Spark pool

Un custom Spark pool lo defines tú con configuración específica. Necesarios cuando:

  • Tienes cargas de trabajo que necesitan más recursos o configuraciones específicas.
  • Necesitas Private Link / Managed Private Endpoint.
  • Quieres fine-tune del cost/performance para un workload concreto.

Puedes configurar:

  • Node family: Memory Optimized (única disponible en Fabric).
  • Node size: desde Small hasta XX-Large.
  • Autoscale: rango min-max de nodos.
  • Dynamic allocation of executors: rango min-max de executors.

4.3 · 🎯 Custom live pools (feature interesante)

Los custom live pools son custom pools con arranque en ~5 segundos (como el starter pool). Requieren un environment con "Full mode for library publishing".

Sin esta configuración, un custom pool tarda ~3 minutos en arrancar. Un starter pool ya está caliente.

4.4 · Prerequisitos para crear custom pools

Para crear un custom pool necesitas:

  1. Rol Admin en el workspace.
  2. Un capacity admin ha habilitado "Customized workspace pools" en Spark Compute settings de la capacidad.

Si no tienes estos permisos → solo puedes usar el starter pool.

4.5 · Cómo se accede a la configuración

Ruta: Workspace → Workspace settingsData Engineering/ScienceSpark settings.

Ahí verás:

  • Tab Pool: gestionar starter pool o crear custom.
  • Tab Environment: entorno por defecto.
  • Tab Compute: configuración de compute.
📏

5. Nodos, tamaños y escalado

5.1 · Tamaños de nodo disponibles

Cada nodo es una VM. Fabric ofrece estos tamaños:

SizevCoresMemoria
Small432 GB
Medium864 GB
Large16128 GB
X-Large32256 GB
XX-Large64512 GB
⚠️ Restricción: X-Large y XX-Large solo se permiten en SKUs no-trial de Fabric.

5.2 · Single-node clusters

Puedes crear pools con 1 solo nodo. En ese caso:

  • Driver y executor comparten la misma máquina.
  • Los recursos se dividen aprox. entre los dos.
  • Ideal para workloads pequeños o desarrollo.

5.3 · Autoscale

Autoscale es el escalado automático del número de nodos según la carga:

  • Configuras un min y max de nodos.
  • Si el pool está inactivo → puede reducir nodos.
  • Si hay picos de trabajo → sube nodos hasta el max.
  • Los nodos se pueden retirar cuando ya no se usan.

Configuración interna:

  • spark.yarn.executor.decommission.enabled = true (por defecto) → apaga nodos infrautilizados agresivamente.
  • Se puede desactivar si prefieres un escalado más suave.

5.4 · Dynamic allocation of executors

Dynamic allocation ajusta el número de executors (no de nodos) según la carga del job:

  • Cuando la app Spark tiene muchas tareas pendientes → pide más executors.
  • Cuando la app está idle o terminó → libera executors.
  • Configuras un min y max de executors.

Es muy útil porque las cargas de un job varían por fase (leer datos → shuffle → agregar → escribir). Los executors necesarios son distintos en cada fase.

5.5 · 🎯 Diferencia autoscale vs dynamic allocation

Esto es sutil pero cae en preguntas:

AutoscaleDynamic allocation
Qué escalaNodos del poolExecutors dentro de la app
Cuándo actúaEntre jobs, según demanda globalDentro de un mismo job
ConfiguraMin/max de nodosMin/max de executors
BeneficioAhorro de costeOptimizar recursos por fase del job

Puedes tener ambos activos → ajuste doble.

5.6 · Right-sizing guidance (recomendaciones oficiales)

De la doc oficial de best practices, según tu patrón de workload:

PatrónRecomendación
Transform-heavy (shuffles, joins)Nodos grandes (16-64 cores)
Bursty o impredecibleAutoscale + Dynamic Allocation
Muchos jobs pequeños paralelosNodos small/medium, runMultiple()
Serial pequeño / desarrolloSmall/medium en single-node
Jobs grandes con particionado conocidoPre-dimensionar manualmente
🏃

6. Runtimes de Spark en Fabric

6.1 · ¿Qué es un runtime?

El runtime determina qué versión de:

  • Apache Spark
  • Delta Lake
  • Python
  • Librerías core

…se usan en tu pool.

6.2 · Runtimes actualmente soportados

Fabric soporta actualmente estos runtimes:

RuntimeSparkDeltaPython
Runtime 1.33.53.13.11
Runtime 2.04.14.13.13+

Fabric NO soporta Spark 3.4 y anteriores. El mínimo es Spark 3.5.

6.3 · Cómo elegir runtime

  • Runtime 1.3 → estable, mucha compatibilidad con librerías Python maduras.
  • Runtime 2.0 → última generación, mejor rendimiento, features nuevas.
⚠️ Autotune solo funciona en Runtime 1.2 (retirado) → no está disponible en runtimes actuales por defecto.

6.4 · Cambiar el runtime

Se configura en:

  • Workspace settings → Spark settings → Environment, o
  • Al crear un environment personalizado.
🎨

7. Environments (entornos personalizados)

7.1 · ¿Qué es un environment?

Un environment es un item de Fabric que empaqueta:

  • Una versión concreta de runtime.
  • Las librerías que quieres instaladas.
  • Configuración de Spark personalizada.
  • Ficheros de recursos que quieres disponibles.
  • Un pool asignado.

Piensa en él como un "perfil personalizado" que puedes reutilizar entre notebooks y Spark jobs. 🎀

7.2 · Por qué usar environments

Casos típicos:

  • Necesitas librerías Python específicas (pandas 2.1.3, scikit-learn 1.4, etc.).
  • Quieres una configuración Spark distinta para ETL grandes vs análisis interactivo.
  • Quieres estandarizar setup entre equipos.
  • Necesitas subir librerías privadas custom.

7.3 · Qué puedes configurar

Al crear un environment puedes:

  • ✅ Especificar el runtime.
  • ✅ Ver las librerías built-in.
  • ✅ Instalar librerías públicas de PyPI.
  • ✅ Instalar librerías custom subiendo un package file.
  • ✅ Asignar un Spark pool.
  • ✅ Configurar propiedades de Spark (override defaults).
  • ✅ Subir ficheros de recursos (JSON de config, modelos ML, etc.).

7.4 · Environment por defecto del workspace

Después de crear un environment, puedes marcarlo como default del workspace. Así, cualquier notebook nuevo lo usará automáticamente.

También puedes anular el default a nivel notebook:

  • Tab del notebook → Workspace default dropdown → seleccionar otro environment.

7.5 · Instalación de librerías en línea (sesión)

Si no quieres crear un environment y solo necesitas una librería para un notebook:

# Instalar una librería en la sesión actual (con %pip)
%pip install pandas==2.1.3

# Con conda
%conda install seaborn=0.13.0
⚠️ La instalación en línea solo dura la sesión. Al terminar, se pierde. Para persistencia → environment.

8. Optimizaciones exclusivas de Fabric

Fabric tiene varias optimizaciones sobre Apache Spark open-source que son específicas de la plataforma. Son muy examinables.

8.1 · 🎯 Native Execution Engine (NEE)

Es un motor de procesamiento vectorizado que ejecuta operaciones Spark directamente sobre la infraestructura del lakehouse.

Beneficios:

  • Rendimiento significativamente mejor para queries sobre datasets grandes en Parquet/Delta.
  • Vectorized processing → procesa muchos valores a la vez.
  • Elimina overhead del engine tradicional.

Cómo activarlo (dos formas):

A) A nivel environment: configurar en el environment:

  • spark.native.enabled = true
  • spark.shuffle.manager = org.apache.spark.shuffle.sort.ColumnarShuffleManager

B) A nivel notebook (per session):

%%configure
{
    "conf": {
        "spark.native.enabled": "true",
        "spark.shuffle.manager": "org.apache.spark.shuffle.sort.ColumnarShuffleManager"
    }
}
⚠️ Ojo cosita: %%configure tiene que ser la primera celda del notebook para que aplique.

8.2 · 🎯 High Concurrency Mode

Qué es: modo que permite compartir una sesión Spark entre múltiples notebooks/usuarios simultáneamente.

Ventajas:

  • Reduce startup overhead (no arranca sesión nueva cada vez).
  • Ahorro de recursos cuando hay muchos users trabajando en paralelo.
  • Mantiene isolation: variables de un notebook no afectan a otro.

Requisitos para compartir sesión:

  • Mismo lakehouse
  • Mismo environment
  • Misma configuración Spark

Cómo activar: Workspace settings → Data Engineering/Science → activar toggle.

También aplicable a Spark jobs (no solo notebooks).

⚠️ Sesiones activas acumulan CU utilization. Timeout por defecto: 20 minutos. Parar sesiones cuando no las uses.

8.3 · 🎯 V-Order

Ya lo vimos en el Módulo 2. Recordamos:

  • Optimización de escritura Parquet exclusiva de Fabric.
  • Activado por defecto (spark.sql.parquet.vorder.enabled = true).
  • Reordena datos para acelerar lecturas desde motores columnar (Power BI, Warehouse).
  • Pequeño overhead al escribir, gran beneficio al leer.

8.4 · Autotune

Autotune es una feature que aprende automáticamente las mejores configuraciones Spark para tus queries recurrentes:

  • Modela las queries que ejecutas.
  • Optimiza settings con cada iteración.
  • Necesita ~20-25 iteraciones para converger.
⚠️ Solo compatible con Runtime 1.2 (retirado). No funciona en runtimes actuales. Si te aparece en preguntas de examen, ten en cuenta que puede estar desactualizado.

8.5 · Automatic MLFlow logging

MLFlow es un tracking system open-source para machine learning. En Fabric:

  • Se usa por defecto para loguear implícitamente actividad de experimentos ML.
  • No necesitas añadir código explícito.
  • Se puede desactivar en workspace settings.

8.6 · Adaptive Query Execution (AQE)

Habilitado por defecto en todos los runtimes de Fabric. Ajusta dinámicamente:

  • Número de shuffle partitions.
  • Estrategias de join.
  • Manejo de skew (particiones desbalanceadas) → spark.sql.adaptive.skewJoin.enabled.
📝

9. Cómo ejecutar código Spark: notebooks vs Spark Job Definitions

Hay dos formas principales de ejecutar código Spark en Fabric. Cae mucho en preguntas escenario.

9.1 · Notebooks

Un notebook es un item interactivo web-based que:

  • Combina celdas de código y celdas de markdown/texto.
  • Ejecuta código cell-by-cell interactivamente.
  • Muestra resultados inline (tablas, gráficos, texto).
  • Se comparte fácilmente.
  • Ideal para exploración, desarrollo, análisis interactivo, ML.

9.2 · Spark Job Definition

Un Spark Job Definition (SJD) es un item que:

  • Ejecuta un script Python/JAR de forma no interactiva.
  • Se puede ejecutar on-demand o por schedule.
  • Se referencia desde pipelines para automatización.
  • Puede especificar ficheros de referencia (código, config).
  • Ideal para producción y ETL automatizado.

9.3 · 🎯 Regla de decisión para el examen

Necesitas…Usa
Explorar datos interactivamenteNotebook
Prototipar transformacionesNotebook
Machine learning con visualizaciones inlineNotebook
Colaborar con markdown y análisisNotebook
ETL programado en producciónSpark Job Definition
Ejecutar un script Python puro sin UISpark Job Definition
Orquestar desde un pipelinePuede ser cualquiera, pero SJD es más "producción-ready"
⚠️ Ojo: notebooks TAMBIÉN se pueden orquestar desde pipelines (actividad "Notebook"). Es muy común en Fabric. Pero para código puramente batch sin UI → SJD es más limpio.

9.4 · Estructura de un notebook

Un notebook tiene:

  • Celdas de código: Python, Spark SQL, Scala, R, HTML, etc.
  • Celdas markdown: documentación, título, imágenes.
  • Un lenguaje primario: el default para todas las celdas nuevas.
  • Un lakehouse asignado por defecto: para que Files/ y Tables/ funcionen sin path completo.

10. Magic commands en notebooks

Los magic commands son comandos especiales que empiezan con % (line magic) o %% (cell magic). Sirven para cambiar comportamiento de una celda o del notebook completo.

10.1 · 🎯 Magics de lenguaje (cambiar idioma de una celda)

MagicLenguajeDescripción
%%pysparkPythonEjecuta Python contra el Spark Context
%%sparkScalaEjecuta Scala contra el Spark Context
%%sqlSpark SQLEjecuta SQL contra el Spark Context
%%sparkrREjecuta R contra el Spark Context
%%htmlHTMLRenderiza HTML
%%csharpC#Ejecuta .NET for Spark C#

Ejemplo:

# Celda en un notebook con lenguaje primario = PySpark

%%sql
SELECT * FROM my_delta_table LIMIT 10

En esa celda concreta, Fabric ejecuta como SQL aunque el notebook sea PySpark. 🌸

10.2 · 🎯 Magics útiles para configuración

MagicUso
%%configureConfigurar Spark (debe ser la 1ª celda)
%pip install XInstalar librería PyPI en la sesión
%conda install XInstalar con conda

Ejemplo de %%configure (activar Native Execution Engine):

%%configure
{
    "conf": {
        "spark.native.enabled": "true",
        "spark.executor.memory": "8g",
        "spark.driver.memory": "4g"
    }
}

10.3 · Magics útiles del kernel IPython

MagicUso
%%timeTiempo total de ejecución de la celda
%%timeitEjecuta múltiples veces y da estadísticas
%%captureCaptura output
%%writefileEscribe el contenido de la celda a un fichero
%%markdownRenderiza como markdown
%%bash / %%shEjecuta comandos de shell

10.4 · ⚠️ Magics soportados en Fabric pipelines

Cuando ejecutas un notebook desde un pipeline (actividad Notebook), solo estos magics están soportados:

  • %%pyspark
  • %%spark
  • %%csharp
  • %%sql
  • %%configure

Si usas otros → pueden fallar en pipeline.

10.5 · Custom magic commands

Puedes definir tus propios magic commands en un notebook y llamarlos desde otro. Útil para funcionalidades reutilizables.

🛠️

11. NotebookUtils (antes MSSparkUtils)

11.1 · ¿Qué es NotebookUtils?

NotebookUtils (antes llamado MSSparkUtils) es una librería built-in para tareas comunes en notebooks:

  • Manejar el file system.
  • Obtener environment variables.
  • Encadenar notebooks.
  • Manejar secrets.
  • Utilities de credentials.

Está disponible en PySpark, Scala, SparkR notebooks y pipelines.

11.2 · ⚠️ Renaming importante

En 2024 se renombró oficialmente:

  • Antes: mssparkutils
  • Ahora: notebookutils

El código antiguo sigue funcionando (backward compatible), pero se recomienda migrar. El namespace mssparkutils se retirará en el futuro.

Requisito: Spark 3.4 (Runtime v1.2) o superior.

11.3 · Uso básico

import notebookutils

# Ver ayuda de todos los comandos disponibles
notebookutils.notebook.help()

11.4 · Ejemplos útiles

Trabajar con file system (fs):

# Listar ficheros
files = notebookutils.fs.ls("Files/my_folder")

# Leer un fichero de texto
content = notebookutils.fs.head("Files/config.json", 1024)  # Primeros 1024 bytes

# Copiar entre ubicaciones
notebookutils.fs.cp("Files/source.csv", "Files/backup/source.csv")

# Borrar
notebookutils.fs.rm("Files/temp.txt")

# Crear directorio
notebookutils.fs.mkdirs("Files/new_folder")

Encadenar notebooks:

# Ejecutar otro notebook y esperar
result = notebookutils.notebook.run("path/to/other_notebook", timeout_seconds=60, args={"param1": "value1"})

# Ejecutar múltiples notebooks en paralelo (DAG)
notebookutils.notebook.runMultiple([
    {"path": "notebook1"},
    {"path": "notebook2", "dependencies": ["notebook1"]},
    {"path": "notebook3", "dependencies": ["notebook1"]}
])

Salir de un notebook con valor:

# Sale del notebook devolviendo un valor al llamador
notebookutils.notebook.exit("Success: 1000 rows processed")

Obtener credentials (secrets):

# Obtener secret de un Azure Key Vault vinculado
secret = notebookutils.credentials.getSecret("kv-name", "secret-name")

11.5 · 🎯 runMultiple() para DAGs

notebookutils.notebook.runMultiple() es SUPER poderoso: ejecuta notebooks con dependencias tipo DAG. Excelente para orquestar pipelines medallion (bronze → silver → gold).

Ejemplo:

DAG = {
    "activities": [
        {"name": "bronze_ingest", "path": "notebooks/01_bronze_ingest"},
        {"name": "silver_transform", "path": "notebooks/02_silver_transform", "dependencies": ["bronze_ingest"]},
        {"name": "gold_aggregate", "path": "notebooks/03_gold_aggregate", "dependencies": ["silver_transform"]}
    ]
}

notebookutils.notebook.runMultiple(DAG)
📊

12. Trabajar con DataFrames

Este es el corazón de PySpark. Vamos a fondo con ejemplos comentados 💕

12.1 · ¿Qué es un DataFrame?

Un DataFrame es una tabla distribuida con filas y columnas. Es lo mismo que un DataFrame de Pandas pero optimizado para procesamiento distribuido.

Nativamente, Spark tiene RDDs (Resilient Distributed Datasets), pero DataFrames son el estándar moderno: más eficientes, con schema, y con API declarativa.

En Fabric usamos DataFrames el 99% del tiempo.

12.2 · Cargar datos en un DataFrame

A) Con schema inferido (opción rápida)

%%pyspark
# Leer un CSV con schema inferido automáticamente por Spark
df = spark.read.load(
    'Files/data/products.csv',   # Ruta al fichero (relativa al lakehouse por defecto)
    format='csv',                # Formato del fichero
    header=True                  # La primera fila son los headers
)

# Mostrar las primeras 10 filas en formato tabla
display(df.limit(10))

Equivalente en Scala:

%%spark
val df = spark.read.format("csv").option("header", "true").load("Files/data/products.csv")
display(df.limit(10))

B) Con schema explícito (mejor rendimiento y validación)

%%pyspark
from pyspark.sql.types import StructType, StructField, IntegerType, StringType, FloatType

# Definir schema explícito
productSchema = StructType([
    StructField("ProductID", IntegerType()),     # Columna entera
    StructField("ProductName", StringType()),    # Columna texto
    StructField("Category", StringType()),
    StructField("ListPrice", FloatType())        # Columna decimal
])

# Cargar CSV con schema (sin header en este ejemplo)
df = spark.read.load(
    'Files/data/product-data.csv',
    format='csv',
    schema=productSchema,     # Se aplica el schema definido
    header=False              # El CSV no tiene fila de cabecera
)

display(df.limit(10))
💡 Tip oficial: "Specifying an explicit schema also improves performance!" Cuando Spark no tiene que inferir el schema, es mucho más rápido.

12.3 · Formatos que puedes leer

Spark en Fabric lee (spark.read.format('X').load()) prácticamente todo:

  • delta (formato preferido en Fabric)
  • parquet
  • csv
  • json
  • orc
  • text
  • xml (con librería extra)
  • avro

Y desde JDBC/ODBC para BBDD.

12.4 · Operaciones típicas: filter, select, groupBy

Seleccionar columnas (select):

# Solo las columnas que necesito
pricelist_df = df.select("ProductID", "ListPrice")

# Sintaxis alternativa (equivalente)
pricelist_df = df["ProductID", "ListPrice"]

Filtrar filas (where o filter):

# Sintaxis con where + notación de columna
bikes_df = df.select("ProductName", "Category", "ListPrice") \
             .where((df["Category"] == "Mountain Bikes") | (df["Category"] == "Road Bikes"))

# Sintaxis con filter (equivalente)
bikes_df = df.filter(df["Category"] == "Mountain Bikes")

# Sintaxis con filter y col() (más "moderna")
from pyspark.sql.functions import col
bikes_df = df.filter((col("Category") == "Mountain Bikes") | (col("Category") == "Road Bikes"))

display(bikes_df)

Agrupar y agregar (groupBy):

# Cuántos productos hay por categoría
counts_df = df.select("ProductID", "Category") \
              .groupBy("Category") \
              .count()

display(counts_df)

Agregar con funciones específicas:

from pyspark.sql.functions import sum, avg, max, min, count

# Estadísticas por categoría
stats_df = df.groupBy("Category").agg(
    count("ProductID").alias("num_products"),
    avg("ListPrice").alias("avg_price"),
    max("ListPrice").alias("max_price"),
    min("ListPrice").alias("min_price"),
    sum("ListPrice").alias("total_price")
)
display(stats_df)

12.5 · 🎯 Method chaining

En PySpark encadenas métodos con . en una sola expresión. Cada método devuelve un nuevo DataFrame (los DataFrames son inmutables).

result = (df.select("ProductName", "Category", "ListPrice")
            .where(col("Category").isin("Mountain Bikes", "Road Bikes"))
            .orderBy("ListPrice", ascending=False)
            .limit(20))

display(result)

Análogo mental: es como componer un pipeline de Power Query paso a paso.

12.6 · Guardar un DataFrame

Guardar como Parquet (formato analítico eficiente):

# Sobrescribir cualquier fichero existente en esa ruta
bikes_df.write.mode("overwrite").parquet('Files/product_data/bikes.parquet')

Modos de escritura (.mode()):

ModoDescripción
overwriteBorra lo existente y escribe de nuevo
appendAñade a lo existente
ignoreNo hace nada si ya existe
error / errorifexists (default)Falla si ya existe

Guardar como Delta (preferido en lakehouse):

# Como tabla Delta gestionada (aparece en el catálogo del lakehouse)
df.write.format("delta").saveAsTable("mi_tabla_delta")

# Como fichero Delta en una ruta específica
df.write.format("delta").mode("overwrite").save("Files/my_delta_folder")
📁

13. Particionado de datos

El particionado es una técnica de optimización que permite a Spark procesar en paralelo mejor y filtrar datos más rápido.

13.1 · Qué es particionar

Particionar = dividir el output en subcarpetas basadas en el valor de una o varias columnas.

Por ejemplo, si tienes datos de ventas y particionas por year y month, obtienes:

Files/sales/
├── year=2024/
│   ├── month=01/
│   │   └── part-00000.parquet
│   ├── month=02/...
│   └── month=12/...
├── year=2025/
│   ├── month=01/...
│   └── ...
└── year=2026/
    └── ...

13.2 · Cómo particionar al escribir

# Particionar por una columna
bikes_df.write.partitionBy("Category").mode("overwrite").parquet("Files/bike_data")

# Particionar por varias columnas (jerarquía)
sales_df.write.partitionBy("year", "month").mode("overwrite").parquet("Files/sales_data")

Las carpetas generadas siguen el formato columna=valor/, que es el estándar Hive-style.

13.3 · Beneficios del particionado

  • Predicate pushdown: si filtras por la columna particionada, Spark lee solo las particiones relevantes, no todo el dataset.
  • 🚀 Paralelismo: cada partición se puede procesar en un executor distinto.
  • 💾 Compresión mejor: datos similares agrupados comprimen mejor.

Ejemplo de filtro que se beneficia:

# Spark solo lee la partición year=2026
df_2026 = spark.read.parquet("Files/sales_data/year=2026")

O con filter automático:

# Spark hace "partition pruning" y solo lee la partición relevante
df = spark.read.parquet("Files/sales_data").filter(col("year") == 2026)

13.4 · ⚠️ Cuándo NO particionar

  • Si tu columna tiene cardinalidad muy alta (miles/millones de valores únicos) → generas miles de particiones minúsculas. Es peor que no particionar.
  • Ejemplos malos: user_id, order_id, timestamp a nivel segundo.
  • Ejemplos buenos: year, country, region, category (poca cardinalidad).

Regla del pulgar: cada partición debería tener ~1 GB aproximadamente. Ni menos ni muchísimo más.

13.5 · Leer datos particionados

Cuando lees, puedes:

A) Leer todo el dataset:

all_data = spark.read.parquet("Files/bike_data")

B) Leer una partición específica (más rápido):

road_bikes = spark.read.parquet("Files/bike_data/Category=Road Bikes")

C) Con wildcards:

# Todas las particiones 2024-2025
some_data = spark.read.parquet("Files/sales_data/year=202[4-5]")
⚠️ Trampa: cuando lees directamente una partición (opción B), la columna particionada desaparece del DataFrame (porque Spark ya sabe que todo son "Road Bikes"). Si haces road_bikes.columns, no verás Category.
🗂️

14. Spark SQL: catalog, temporary views y tablas

Spark tiene un catalog (metastore) que registra objetos relacionales para consultarlos con SQL.

14.1 · Catalog

El Spark catalog es el metastore que registra:

  • Databases (schemas)
  • Tables (managed o external)
  • Views (persistentes o temporary)
  • Functions

En Fabric, el catalog del lakehouse tiene tus tablas Delta bajo dbo (o schemas custom si tienes schema-enabled lakehouse).

14.2 · Temporary views

Una temporary view es una vista efímera que existe solo durante la sesión. Se autodestruye al cerrar la sesión.

Uso típico: hacer queriable un DataFrame con SQL sin persistirlo como tabla.

# Convierto un DataFrame en una vista temporal consultable
df.createOrReplaceTempView("products_view")

# Ahora puedo consultar con SQL
result = spark.sql("SELECT * FROM products_view WHERE Category = 'Bikes'")
display(result)

# O directamente con magic
%%sql
SELECT * FROM products_view WHERE Category = 'Bikes'

Al terminar la sesión → la vista desaparece. Los datos originales del DataFrame también desaparecen si no los guardaste.

14.3 · Tablas managed

Las tablas managed son las que Fabric gestiona completamente (metadata + datos). Los datos van a Tables/ del lakehouse por defecto.

Crear con saveAsTable:

# Convierte el DataFrame en una tabla Delta gestionada
df.write.format("delta").saveAsTable("products")

# Con schema específico
df.write.format("delta").saveAsTable("sales.orders")  # requiere schemas habilitados

Crear vacía con spark.catalog:

from pyspark.sql.types import *

schema = StructType([
    StructField("id", IntegerType()),
    StructField("name", StringType())
])

spark.catalog.createTable("empty_table", schema=schema, source="delta")
⚠️ Al borrar una tabla managed, se borran también los datos.

14.4 · Tablas external (matices ya vistos)

Como vimos en el Módulo 2, en lakehouses de Fabric la recomendación es usar shortcuts, no external tables. Pero técnicamente:

# External table: apunta a datos en otra ubicación
df.write.format("delta").saveAsTable(
    "external_products",
    path="Files/external/products"
)

Al borrar la tabla → solo se borra el metadata, los datos permanecen.

API alternativa:

spark.catalog.createExternalTable("external_products", "Files/external/products", source="delta")

14.5 · 🎯 Formato preferido en Fabric

Delta es el formato preferido en Fabric para tablas:

  • Transacciones ACID
  • Time travel
  • Schema enforcement
  • V-Order
  • Aparece en SQL analytics endpoint

Otros formatos (Parquet, CSV) NO aparecen en SQL analytics endpoint.

14.6 · Consultar con Spark SQL

Con la API spark.sql():

# Devuelve un DataFrame
bikes_df = spark.sql("""
    SELECT ProductID, ProductName, ListPrice
    FROM products
    WHERE Category IN ('Mountain Bikes', 'Road Bikes')
""")

display(bikes_df)

Con magic %%sql:

%%sql
SELECT Category, COUNT(ProductID) AS ProductCount
FROM products
GROUP BY Category
ORDER BY Category

Los resultados se muestran automáticamente como tabla en el notebook.

14.7 · 🎯 Casos de uso

  • spark.sql(): cuando quieres usar el resultado en Python (asignar a variable, seguir procesando).
  • %%sql: cuando quieres solo ver el resultado en pantalla, exploración rápida.

14.8 · Referencias four-part namespace

Si tienes schemas habilitados:

-- Referencia completa
SELECT * FROM my_workspace.my_lakehouse.sales.orders

-- Referencia dentro del mismo workspace
SELECT * FROM my_lakehouse.sales.orders

-- Referencia dentro del mismo lakehouse (schema explícito)
SELECT * FROM sales.orders

-- Referencia al schema por defecto (dbo)
SELECT * FROM orders
📈

15. Visualización en notebooks

15.1 · Charts built-in del notebook

Cuando ejecutas una celda que devuelve un DataFrame o SQL query, el resultado sale como tabla. Puedes cambiar a modo chart desde la UI del notebook:

  • Barra debajo de los resultados → botón Chart.
  • Configurar chart type, dimensiones, agregaciones desde la UI.

Es rápido para exploración pero limitado en personalización.

15.2 · Matplotlib

Matplotlib es la librería base de visualización en Python. Está pre-instalada en el runtime de Fabric.

Requisito importante: Matplotlib necesita datos en Pandas DataFrame, no en Spark DataFrame.

Para convertir:

# Convertir de Spark a Pandas (¡cuidado con datasets grandes!)
pandas_df = spark_df.toPandas()
⚠️ Ojo cosita: .toPandas() trae todos los datos al driver (a memoria local). Si el dataset es enorme → out of memory. Solo usa .toPandas() con datos ya agregados o pequeños.

Ejemplo completo:

from matplotlib import pyplot as plt

# 1. Agregar los datos con Spark (procesamiento distribuido)
data = spark.sql("""
    SELECT Category, COUNT(ProductID) AS ProductCount
    FROM products
    GROUP BY Category
    ORDER BY Category
""").toPandas()   # <-- Convertimos a Pandas DESPUÉS de agregar

# 2. Limpiar cualquier plot previo
plt.clf()

# 3. Crear la figura con tamaño personalizado
fig = plt.figure(figsize=(12, 8))

# 4. Crear el bar chart
plt.bar(x=data['Category'], height=data['ProductCount'], color='orange')

# 5. Personalizar
plt.title('Product Counts by Category')
plt.xlabel('Category')
plt.ylabel('Products')
plt.grid(color='#95a5a6', linestyle='--', linewidth=2, axis='y', alpha=0.7)
plt.xticks(rotation=70)   # Rotar etiquetas del eje X

# 6. Mostrar
plt.show()

15.3 · Seaborn (más bonito)

Seaborn está construido sobre Matplotlib pero con estilos más profesionales y bonitos:

import seaborn as sns
import matplotlib.pyplot as plt

# Setup del estilo
sns.set_theme(style="whitegrid")

# Convertir a Pandas
data = spark.sql("SELECT * FROM sales").toPandas()

# Ejemplo: scatter plot con color por categoría
sns.scatterplot(data=data, x="ListPrice", y="Sales", hue="Category")
plt.title("Precio vs Ventas por Categoría")
plt.show()

15.4 · Otras librerías útiles

Todas se instalan con %pip install X o vía environment:

  • Plotly: gráficos interactivos.
  • Bokeh: gráficos interactivos web-based.
  • Altair: gráficos declarativos.

15.5 · 🎯 Best practice

Regla del pulgar:

  1. Haz agregación en Spark (distribuido, escalable).
  2. Convierte a Pandas solo el resultado agregado (pequeño).
  3. Visualiza con Matplotlib/Seaborn/Plotly.

NO al revés: no traigas todo a Pandas y luego agregues. Perderás la ventaja de Spark.

🎛️

16. Elección de kernel: Python vs Spark

Los notebooks de Fabric soportan 3 kernels:

  • Python (kernel puro Python, sin Spark)
  • Spark (con soporte PySpark, SparkSQL, Scala, SparkR)
  • T-SQL (T-SQL puro)

Vamos a comparar Python vs Spark que es donde suele haber dudas.

16.1 · Python kernel

  • Ejecuta Python puro (Pandas, NumPy, scikit-learn, etc.).
  • Single-node: todo corre en una máquina, sin distribución.
  • Startup más rápido.
  • Ideal para: datasets pequeños-medianos, ML clásico, análisis exploratorio pequeño, scripts simples.
  • NO puede acceder a features distribuidas de Spark.

16.2 · Spark kernel

  • Ejecuta PySpark, SparkSQL, Scala, SparkR — todos con el mismo compute.
  • Multi-node: procesamiento distribuido.
  • Startup más lento (a menos que uses starter pool o live pool).
  • Ideal para: datasets grandes, ETL, Delta Lake, transformaciones distribuidas.
  • Puede acceder al catalog del lakehouse, Delta features, escritura a Tables.

16.3 · 🎯 Consideraciones para el examen

Microsoft es explícito en la doc: "The choice of notebook kernel isn't simply about cost or data size."

Factores importantes:

FactorA favor de Python kernelA favor de Spark kernel
Volumen de datosPocos MB - GBCientos de GB - TB
Crecimiento futuroEstableVa a crecer mucho
Uso de Delta Lake featuresBajoAlto (schema evolution, time travel, etc.)
ConcurrenciaBajaAlta
Librerías PythonSí (mejor compatibilidad)Sí, con matices
CostMenor overhead por sesiónMayor overhead pero escalable

16.4 · Reglas prácticas

  • Un notebook de exploración con pandas y matplotlib sobre 100 MB → Python kernel.
  • Un notebook de ETL bronze→silver con 500 GB de datos → Spark kernel.
  • Un notebook ML con 5 GB de features y XGBoost → puede ir en cualquiera, pero Python kernel + Pandas suele ser más simple.
  • Cualquier operación que escribe tablas Delta gestionadas del lakehouseSpark kernel (Python puro no tiene acceso al catalog).
⚠️

17. Trampas típicas y confusiones frecuentes

Trampa 1: "El pool y el environment son lo mismo"

Distintos. El pool es el compute (nodos, memoria). El environment es la config (runtime, libraries, Spark props). Un environment usa un pool.

Trampa 2: "Autoscale = Dynamic allocation"

❌ Autoscale escala nodos. Dynamic allocation escala executors. Son independientes y se pueden combinar.

Trampa 3: "El starter pool arranca al instante"

✅ Sí (~segundos), pero custom pools sin live mode tardan ~3 min. Custom live pools ~5 seg.

Trampa 4: "V-Order lo tengo que activar yo"

❌ Está activado por defecto en Fabric. Solo lo desactivas si tienes un motivo específico.

Trampa 5: "Native Execution Engine también está por defecto"

NO. Hay que activarlo explícitamente con %%configure o en el environment.

Trampa 6: "En Fabric hay Spark 3.4"

Fabric NO soporta Spark 3.4 ni anteriores. Runtimes actuales: 1.3 (Spark 3.5) y 2.0 (Spark 4.1).

Trampa 7: "Autotune está siempre disponible"

❌ Autotune solo funciona en Runtime 1.2 (retirado). No está en runtimes actuales.

Trampa 8: "%%sql corre en el SQL analytics endpoint"

NO. %%sql en un notebook corre en el Spark SQL engine del pool, no en el SQL analytics endpoint. Son motores distintos.

Trampa 9: ".toPandas() es siempre seguro"

❌ Trae TODOS los datos al driver. Con datasets grandes → out of memory. Úsalo solo con datos ya agregados o pequeños.

Trampa 10: "Puedo particionar por cualquier columna"

❌ Particionar por columnas de alta cardinalidad genera miles de ficheros pequeños → peor performance. Usa columnas de baja cardinalidad.

Trampa 11: "Cuando leo una partición específica, la columna sigue en el DataFrame"

❌ La columna particionada desaparece del DataFrame cuando lees directamente una partición: spark.read.parquet("Files/data/Category=Road Bikes") → no verás la columna Category.

Trampa 12: "Todos los magics funcionan en pipelines"

En pipelines solo funcionan: %%pyspark, %%spark, %%csharp, %%sql, %%configure. Otros fallan.

Trampa 13: "MSSparkUtils y NotebookUtils son distintos"

❌ Son el mismo paquete renombrado. mssparkutils seguirá funcionando por compatibilidad, pero notebookutils es el nuevo nombre oficial.

Trampa 14: "High concurrency mode se activa por defecto"

❌ Hay que activarlo manualmente en workspace settings.

🆚

18. Comparativas clave (chuleta)

Starter pool vs Custom pool

Starter poolCustom pool
CreaciónAutomático por workspaceManual
Arranque~segundos~3 min (o ~5 seg si es live)
ConfigPredefinida (medium)Total control
Private Link/Endpoint
CostsOptimizadoDepende de config

Notebook vs Spark Job Definition

NotebookSpark Job Definition
Interactividad✅ Cell-by-cell❌ Script completo
Visualización inline
Markdown/docs
Programación desde pipeline
Ideal paraDesarrollo, exploración, MLProducción, batch ETL

Autoscale vs Dynamic allocation

AutoscaleDynamic allocation
EscalaNodos del poolExecutors dentro del job
Cuándo actúaEntre jobs / globalDentro de un mismo job
ConfigMin/max nodosMin/max executors

Python kernel vs Spark kernel

PythonSpark
Distribuido
StartupRápidoLento (a menos que live)
Delta lakehouseLimitadoCompleto
Ideal paraDatasets pequeños, ML clásicoETL grandes, Delta

Magic commands útiles

MagicPara qué
%%pysparkCambiar celda a PySpark
%%sqlCambiar celda a SparkSQL
%%sparkCambiar celda a Scala
%%configureConfigurar Spark (1ª celda)
%pip install XInstalar librería en sesión
%%timeMedir tiempo de la celda

Optimizaciones de Fabric

FeatureActivada por defectoCómo activar
V-Order(ya activa)
Adaptive Query Execution(ya activa)
Native Execution Engine%%configure o environment
High Concurrency ModeWorkspace settings
Autotune⚠️Solo Runtime 1.2 (retirado)
🚀 ¡Módulo 3 completado! Ya dominas el motor Spark de Fabric. Sigue con el Módulo 4: Tablas Delta Lake. 🌸