Usar eventstreams en tiempo real en Microsoft Fabric
La herramienta principal de ingesta y transformación de streaming: sources, destinations, transformaciones no-code, SQL operator, las 5 windowing functions, derived streams y DeltaFlow. 🌸
IntermedioAvanzado
🎯
1. Objetivos y encaje en el DP-700
¿Qué te enseña este módulo?
Este módulo es donde profundizamos en los Eventstreams — la herramienta principal para la ingesta y transformación de datos en streaming en Fabric. Los objetivos oficiales:
Establecer sources y destinations en los Fabric Eventstreams.
Capturar, transformar y enrutar datos usando Fabric Eventstreams.
Peso en el examen DP-700
Los Eventstreams aparecen con fuerza en el Dominio 2 (Ingest and transform) y también en el Dominio 3 (troubleshooting):
Dominio
Cómo aparecen
Ingest and transform (30-35%)
Elección del motor de streaming, transformaciones, windowing functions, routing, content-based routing
Monitor and optimize (30-35%)
Troubleshooting de Eventstreams, workspace monitoring, rendimiento
🎯 Qué es CRÍTICO dominar:
Los 5 tipos de windowing functions (esto CAE en el examen fijo)
Derived streams para content-based routing
Elegir entre Eventstream y Spark Structured Streaming
Cuándo usar el SQL code operator vs las transformaciones no-code
Todos los destinations y sus casos de uso
DeltaFlow para CDC (feature moderna)
🌸 Sobre las windowing functions: son la piedra angular del procesamiento de streams y aparecen en muchísimas preguntas de examen. Te las explico con muchos ejemplos y diagramas mentales para que se te queden bien pegadas 💕
🧩
2. Anatomía de un Eventstream
¿Qué es un Eventstream?
Un Eventstream es una feature de Fabric que crea un pipeline que:
Ingesta eventos de fuentes de datos en streaming.
Los procesa con transformaciones opcionales.
Los entrega a varios destinos.
Es el mecanismo de entrega que lleva los eventos desde donde ocurren hasta donde necesitan procesarse.
El canvas visual
Los eventstreams se diseñan en un canvas visual: drag and drop de nodos (sources, transformaciones, destinations), viendo los datos de eventos fluyendo en tiempo real. Sin código y sin gestionar infraestructura.
Los 3 componentes obligatorios
Sources 📥 — de dónde vienen los eventos.
Transformations 🔄 — el procesamiento opcional que se aplica.
Destinations 📤 — dónde van los eventos procesados.
Ejemplo mental
Imagina que tienes datos de bicis compartidas:
Bicycles (source) ← datos raw de alquiler de bicis
↓
GroupByStreet ← transform: suma bicis por calle
↓
Bikes-by-street-table ← destination: tabla KQL
Todo esto en el canvas visual, drag-and-drop, sin código.
🎨
3. Enhanced vs Standard capabilities
Dos versiones de Eventstream
A) Standard capabilities — la versión original, con interfaz básica y sources y destinations más limitados.
B) Enhanced capabilities ⭐ — la versión moderna (RECOMENDADA): más sources y destinations, mejor UI/UX y nuevas features (SQL operator, DeltaFlow, etc.).
📌 Cuál elegir:Enhanced siempre, a menos que estés manteniendo un eventstream legacy. Al crear un eventstream nuevo en 2026, la opción por defecto es Enhanced. Este manual se centra en Enhanced.
📥
4. Sources: catálogo completo
Los eventstreams pueden ingestar de muchos tipos de fuentes:
Microsoft sources
Azure Event Hubs ⭐ (el más común)
Azure IoT Hub
Azure Service Bus
Feeds de Change Data Capture (CDC) en servicios de base de datos: Azure SQL Database CDC, PostgreSQL CDC, MySQL CDC, SQL Server on VM CDC y Azure SQL Managed Instance CDC.
Azure events
Eventos de Azure Blob Storage (BlobCreated, BlobDeleted).
Fabric events
Cambios en items del workspace (creación, modificación, eliminación).
Cambios de datos en los data stores de OneLake.
Eventos de jobs de Fabric (pipeline start/succeed/fail, etc.).
Fabric incluye streams de muestra pregenerados para practicar: Bicycles (alquiler de bicis), Yellow Taxi (taxis de NYC) y Stock Market (precios de acciones simulados). Ideales para aprender sin montar infraestructura real.
Business events (Preview)
Los Fabric business events (feature más nueva) permiten publicar business events desde apps con schemas definidos y reaccionar a esos eventos con Activator + User Data Functions. Es el patrón moderno de arquitectura orientada a eventos.
🎯 Elegir el source correcto
Escenario
Source
Dispositivos IoT enviando telemetría
Azure IoT Hub
Clúster Kafka on-prem
Apache Kafka
BBDD transaccional con CDC
CDC feed
Ficheros que aterrizan en Blob
Eventos de Azure Blob Storage
App custom generando eventos
Custom app source
Multi-cloud (GCP, AWS)
Google Cloud Pub/Sub, Kinesis
⚙️
5. Cómo configurar sources
Dos formas
A) Crear un source nuevo desde el canvas del eventstream: click en "Add source", elegir el tipo y configurar la conexión y las credenciales.
B) Conectar a un source existente desde el Real-Time Hub: el Real-Time Hub tiene un catálogo de sources ya configurados; te suscribes desde el hub y es reutilizable entre eventstreams.
Datos de configuración típicos
Para Azure Event Hubs: connection string o Shared Access Signature, consumer group (para no interferir con otros consumidores) y content format (JSON, Avro, etc.).
Para Kafka: bootstrap servers, nombre del topic y autenticación (SASL, SSL, etc.).
Para CDC: conexión a la base de datos, tablas a monitorizar y modo de output (analytics-ready recomendado).
Un eventstream, varios sources
Puedes conectar varios sources a un mismo eventstream. Ideal para hacer union de streams similares, join entre streams distintos o enriquecimiento cross-source.
📤
6. Destinations: catálogo completo
Los destinos son endpoints donde tus datos procesados quedan disponibles para queries, reports, dashboards, alertas y acciones.
Eventhouse
Ingestar tus datos de eventos en tiempo real en un Eventhouse para consultarlos con KQL.
Cuándo usarlo: analítica en tiempo real con KQL, análisis de series temporales, detección de anomalías e ingesta de alta velocidad.
Lakehouse
Transformar los eventos e ingestarlos en el lakehouse en formato Delta Lake.
Cuándo usarlo: persistir datos históricos, consumirlos desde Power BI Direct Lake, análisis batch posterior sobre datos de streaming y data engineering híbrido batch+streaming.
Derived stream ⭐
Versión transformada del stream de datos original que habilita el content-based routing. Permiten filtrar y transformar una vez, enrutar subconjuntos según el contenido y reutilizarlos desde varios sistemas downstream.
Ejemplo: filtrar datos de sensores IoT → alertas de temperatura alta a Activator, medias horarias a una KQL database. Los vemos en detalle en la sección de content-based routing.
Fabric Activator
Conectar directamente los datos de eventos en tiempo real a un motor de detección de eventos que dispara acciones automáticamente cuando detecta patrones.
Acciones posibles: enviar notificaciones (email, Teams), lanzar workflows de Power Automate, ejecutar pipelines de Fabric y ejecutar notebooks de Fabric.
Custom endpoint
Enrutar los eventos en tiempo real a un endpoint externo custom. Cuándo usarlo: integraciones con sistemas fuera de Fabric, aplicaciones custom y servicios de terceros.
🎯 Varios destinations en el mismo eventstream: puedes conectar a varios destinos simultáneamente sin que se afecten ni colisionen entre ellos.
Un source, un procesado, 4 destinations — todo en paralelo.
⚙️
7. Cómo configurar destinations
En el canvas
Después de conectar un source (y opcionalmente transformaciones), puedes añadir un destino:
Click en Add destination.
Elegir el tipo.
Configurar el destino (workspace, item, tabla, etc.).
Mapear columnas si aplica.
Configuración típica por destino
Eventhouse destination: destino = workspace + Eventhouse + KQL database + tabla; modo de ingesta directo o vía update policy; y mapeo de schema.
Lakehouse destination: destino = workspace + Lakehouse + tabla; settings de la tabla Delta; y estrategia de particionado.
Activator destination: seleccionas el Activator existente o creas uno nuevo. Los eventos fluyen y Activator los evalúa contra sus rules.
Custom endpoint: URL del endpoint, autenticación (Bearer token, connection string) y formato (JSON, Avro).
🔄
8. Transformaciones no-code (7 tipos)
Eventstream ofrece 7 transformaciones no-code que arrastras al canvas. Son las herramientas del día a día.
Filter 🔍
Filtra eventos según el valor de un campo. Ejemplos de condiciones: temperature > 80°, status = "error", customer_type = "premium". Uso típico: quitar eventos inválidos o irrelevantes antes del almacenamiento.
Manage fields ✏️
Permite añadir campos calculados (derivados de otros), quitar campos innecesarios, cambiar tipos de datos y renombrar campos. Uso típico: preparar el schema para el destino.
Ejemplo: partiendo de [temp_celsius, timestamp_utc, device_id] → añadir temp_fahrenheit = temp_celsius * 1.8 + 32, renombrar device_id → deviceId y cambiar el tipo de timestamp_utc → datetime.
Aggregate 📊
Calcula una agregación (Sum, Min, Max, Average) cada vez que ocurre un evento nuevo sobre un periodo de tiempo. Puedes tener varias agregaciones en la misma transformación, renombrar las columnas calculadas y filtrar la agregación por dimensiones.
Group by 🌸 (CRÍTICO — contiene el windowing)
Calcula agregaciones entre eventos dentro de ventanas temporales. Ejemplos: totales de ventas por hora, medias diarias de temperatura.
Soporta varias ventanas temporales: Tumbling (intervalos fijos), Hopping (intervalos solapados) y Sliding (solapados cuando cambia el contenido). Los detalles los vemos en la sección específica.
Union 🔗
Conecta 2 o más nodos del canvas y añade los eventos con campos compartidos (mismo nombre y tipo de dato) a una tabla combinada.
Reglas: los campos que no coinciden se descartan del output. Similar al UNION de SQL pero para streams. Uso típico: consolidar streams de varios sources con schema parecido.
Ejemplo: stream 1 con sensores del edificio A [device_id, temp, timestamp] y stream 2 del edificio B [device_id, temp, timestamp, humidity] → el union da [device_id, temp, timestamp] (humidity se cae).
Join 🔀
Combina datos de 2 streams según una condición de coincidencia. Similar al JOIN de SQL (INNER JOIN, LEFT JOIN, etc.). Uso típico: enriquecimiento cross-stream (correlacionar eventos).
Ejemplo: stream A con orders [order_id, customer_id, amount] y stream B con customers [customer_id, name, region], join por customer_id → orders enriquecidas con la info del cliente.
Expand 📤
Transformación de arrays que crea una fila nueva por cada valor dentro de un array. Similar a explode() en Spark. Uso típico: normalizar eventos con arrays anidados.
Para equipos code-first, Eventstream incluye un operador SQL donde escribes expresiones SQL directamente. Soporta windowing, agregaciones, joins, varias queries en el mismo operador y content-based routing (misma query, distintos outputs).
Cuándo usarlo
SQL operator ✅ cuando tu equipo viene de Stream Analytics (misma sintaxis), la lógica es compleja y necesitaría muchas transformaciones no-code encadenadas, quieres content-based routing en una sola query, o prefieres código declarativo a lo visual.
Transformaciones no-code ✅ cuando eres analista low-code, la lógica es simple y clara, o el equipo quiere auditabilidad visual.
Ejemplo
SELECT
device_id,
AVG(temperature) AS avg_temp,
COUNT(*) AS event_count,
System.Timestamp() AS window_end
FROM
IoT_Stream
WHERE
temperature IS NOT NULL
GROUP BY
device_id,
TumblingWindow(minute, 5)
Se traduce como: agrupa los eventos por dispositivo y ventanas de 5 minutos, y calcula la media y el conteo.
Content-based routing en SQL
Puedes tener varios SELECT en un operador SQL, cada uno con un output distinto:
-- Output 1: alertas de temperatura alta
SELECT * INTO high_temps
FROM IoT_Stream
WHERE temperature > 80;
-- Output 2: lecturas normales
SELECT * INTO normal_readings
FROM IoT_Stream
WHERE temperature <= 80;
Esto permite hacer routing a distintos destinos desde un solo operador SQL.
⭐
10. Windowing functions: el corazón del streaming
Este es EL tema más examinable del módulo. Presta atención especial 🌸
¿Qué son las windowing functions?
En streaming, una window es un intervalo de tiempo sobre el que agrupas eventos para hacer cálculos agregados.
Por qué son necesarias: los streams son infinitos. Un GROUP BY clásico requeriría procesar todos los eventos, algo imposible con streams. Las windows te permiten agrupar eventos en intervalos temporales limitados.
Los 5 tipos de windows
Tumbling Window — fijas, sin solapamiento
Hopping Window — fijas, con solapamiento (hops)
Sliding Window — móvil, disparada por eventos
Session Window — grupo de eventos cercanos con timeout
Snapshot Window — mismo timestamp
Fabric Eventstream soporta las 5 (algunas también en Azure Stream Analytics).
🎯 Regla clave: output al final.Todas las operaciones de windowing emiten eventos al FINAL de la ventana. Cada window tiene un start time, un end time y un output: un único evento por ventana con la agregación calculada, con el timestamp del end time.
Sintaxis en el operador SQL
Se usan en la cláusula GROUP BY:
SELECT ... FROM stream
GROUP BY <campos>, <WindowFunction>(<params>)
Límites generales
Tamaño máximo de ventana: 7 días para todos los tipos.
Puedes agregar sobre varias ventanas en el mismo GROUP BY usando la función Windows().
Vamos a ver cada tipo en detalle 🌸
🌸
11. Tumbling Window a fondo
Definición
Las tumbling windows son una serie de intervalos de tamaño fijo, sin solapamiento y contiguos. Piénsalos como cubos alineados en el tiempo, todos del mismo tamaño y sin solaparse.
🌷 Analogía kawaii: como magdalenas en una bandeja: cada una en su hueco, del mismo tamaño y sin tocarse. Cada evento cae en exactamente una magdalena (una ventana).
SELECT
TollId,
COUNT(*) AS car_count
FROM Input TIMESTAMP BY EntryTime
GROUP BY
TollId,
TumblingWindow(minute, 5)
Traducción: agrupa los eventos por TollId y ventanas tumbling de 5 minutos, cuenta los coches por ventana y emite el output cada 5 minutos con el conteo final.
Comportamiento de inclusión
Por defecto, incluye el final de la ventana (12:05) y excluye el principio (12:00). Ejemplo: la ventana 12:00-12:05 incluye los eventos exactamente a las 12:05, pero no los de las 12:00 (esos van a la ventana 11:55-12:00). El offset cambia este comportamiento si lo necesitas.
Cuándo usarla
Informes horarios (TumblingWindow(hour, 1)).
Rollups de 5 minutos para dashboards.
Agregaciones diarias (TumblingWindow(day, 1)).
Cualquier caso donde quieras cubos temporales limpios sin solapamiento.
🌸
12. Hopping Window a fondo
Definición
Las hopping windows son ventanas programadas con solapamiento. Como las tumbling, pero se solapan entre ellas — van "saltando" (hop) hacia delante en el tiempo por un periodo fijo.
🌷 Analogía kawaii: como los pasos de un baile con solapamiento: cada paso empieza antes de que termine el anterior. Un evento puede estar en varios pasos (ventanas) a la vez.
Los 3 parámetros
Requiere 3 parámetros (a diferencia de la tumbling, que solo necesita 1):
timeunit: unidad de tiempo (minute, second, etc.).
windowsize: cuánto dura cada ventana.
hopsize: cuánto avanza la ventana respecto a la anterior.
🎯 Regla clave: una tumbling window es un caso especial de hopping window donde hop = size (sin solapamiento). Si haces HoppingWindow(minute, 5, 5) equivale a TumblingWindow(minute, 5).
Ejemplo con solapamiento
HoppingWindow(minute, 10, 5): duración de 10 minutos y hop de 5 minutos (avanza 5 min entre ventanas). Ventanas generadas:
Un evento a las 10:07 estaría en la Window 1 Y en la Window 2 a la vez.
Ejemplo SQL completo
SELECT
Region,
AVG(Temperature) AS avg_temp
FROM SensorInput TIMESTAMP BY EventTime
GROUP BY
Region,
HoppingWindow(minute, 10, 5)
Interpretación: media de temperatura por región, con ventanas de 10 minutos que "saltan" cada 5 minutos. Output cada 5 minutos con la media de los últimos 10 minutos.
Cuándo usarla
Medias móviles con actualizaciones más frecuentes que la duración del cálculo.
Dashboards suaves que quieren actualizarse a menudo con métricas de ventanas amplias.
Ejemplo: un KPI de dashboard actualizado cada minuto con datos de los últimos 15 minutos → HoppingWindow(minute, 15, 1).
🌸
13. Sliding Window a fondo
Definición
Las sliding windows solo emiten output cuando el contenido de la ventana realmente cambia, es decir, cuando un evento entra o sale de la ventana.
Diferencia crítica con tumbling/hopping: las tumbling y hopping emiten a intervalos fijos; la sliding emite de forma event-driven (cuando algo cambia).
🌷 Analogía kawaii: como un detector de movimiento: solo emite señal cuando algo entra o sale de su rango. Si nada se mueve, no emite nada.
Características
✅ Cada ventana tiene al menos 1 evento (garantizado).
✅ Los eventos pueden pertenecer a varias sliding windows.
✅ Output on-demand, no programado.
📏 Solo necesita 1 parámetro: la duración.
Sintaxis
SLIDINGWINDOW(timeunit, windowsize)
Simple — no hay hopsize ni offset.
Ejemplo
Input:
Stamp
CreatedAt
Topic
1
2021-10-26T10:15:10
Streaming
5
2021-10-26T10:15:12
Streaming
9
2021-10-26T10:15:15
Streaming
7
2021-10-26T10:15:15
Streaming
8
2021-10-26T10:15:27
Streaming
Query:
SELECT
System.Timestamp() AS WindowEndTime,
Topic,
COUNT(*) AS Count
FROM TwitterStream TIMESTAMP BY CreatedAt
GROUP BY Topic, SlidingWindow(second, 10)
HAVING COUNT(*) >= 3
Output (solo cuando cambia el conteo Y cumple el HAVING >= 3):
WindowEndTime
Topic
Count
2021-10-26T10:15:15
Streaming
4
2021-10-26T10:15:20
Streaming
3
Interpretación: solo emite cuando la ventana de 10 segundos tiene 3 o más eventos.
Cuándo usarla
Alertas (alertar solo cuando cambia la condición).
Detección de patrones que ocurren con timing irregular.
Casos donde quieres evitar emitir eventos "sin cambio".
🌸
14. Session Window a fondo
Definición
Las session windows agrupan eventos que llegan en momentos cercanos, filtrando los periodos en los que no hay datos. Cada "sesión" de actividad se agrupa como una ventana, y cuando hay un hueco de inactividad, la sesión termina.
🌷 Analogía kawaii: como las sesiones de conversación en WhatsApp: se agrupan los mensajes seguiditos como una charla, pero cuando pasa un rato sin mensajes, empieza una sesión nueva.
Los 3 parámetros
Timeout: cuánto esperar sin datos nuevos antes de cerrar la sesión.
Maximum duration: el tiempo máximo que puede durar una sesión (tope de seguridad).
Partition (opcional): partition key para agrupar sesiones por entidad.
Comportamiento
Si llega un evento → empieza (o continúa) una sesión.
Si en timeout segundos NO llega otro evento → la sesión se cierra.
Si la sesión alcanza el maximum_duration → se cierra forzosamente (aunque sigan llegando eventos).
Con partition, se agrupan los eventos solo dentro de la misma key.
Casos de uso ideales
Seguimiento de sesiones de usuario: agrupar las acciones de un usuario en su sesión web.
Actividad de dispositivos: periodos de actividad de un dispositivo IoT.
Detección de anomalías: periodos de errores agrupados.
Ejemplo mental
Un usuario navegando un e-commerce: 10:00 visita la home, 10:01 hace click en el producto A, 10:02 lo añade al carrito, (hueco de 30 min sin actividad), 10:32 vuelve y hace checkout.
Con una SessionWindow con timeout de 20 minutos: Sesión 1 = 10:00-10:02 (3 eventos) y Sesión 2 = 10:32 (1 evento).
🌸
15. Snapshot Window a fondo
Definición
Las snapshot windows agrupan los eventos que tienen el mismo timestamp. Es la más simple: sin parámetros, agrupa por el timestamp exacto del sistema.
Características
❌ No requiere parámetros — usa el system time.
✅ Ideal para agrupar eventos exactamente simultáneos.
🎯 Muy poco común, pero útil en escenarios específicos.
Cuándo usarla
Cuando tienes varios eventos con exactamente el mismo timestamp que quieres agrupar.
Sistemas donde el system time es la única "agrupación" válida.
🎯
16. Cómo elegir el window correcto
Árbol de decisión
Necesito...
Window
Intervalos fijos sin solapamiento
Tumbling
Intervalos fijos con solapamiento
Hopping
Emitir solo cuando cambia el contenido
Sliding
Agrupar por "sesión" de actividad
Session
Agrupar eventos con el mismo timestamp
Snapshot
Preguntas guía
P1: ¿quiero solapamiento entre ventanas? Sí → Hopping o Sliding. No → Tumbling o Session.
P2: ¿quiero output a intervalos regulares o event-driven? Regular → Tumbling o Hopping. Event-driven → Sliding o Session.
P3: ¿los eventos forman "sesiones" con periodos de inactividad? Sí → Session.
P4: ¿quiero agrupar eventos por timestamp exacto? Sí → Snapshot.
Comparativa visual
Tumbling: [....][....][....][....] cubos separados
Hopping: [....] solapamiento
[....]
[....]
[....]
Sliding: |------| emite cuando entra/sale un evento
|------|
Session: [.......] gap [....] gap [........]
(timeout) (timeout)
Snapshot: [| | ||| | |] agrupa los mismos timestamps
Ejemplos reales por window
Window
Ejemplo real
Tumbling
"Ventas totales cada hora"
Hopping
"Media móvil de temperatura cada minuto sobre los últimos 15 min"
Sliding
"Alertar cuando en los últimos 10 seg hay 3 o más errores"
Session
"Sesiones de navegación del usuario con 20 min de timeout"
Snapshot
"Todas las cotizaciones bursátiles del mismo tick"
🔀
17. Content-based routing con Derived streams
¿Qué es el content-based routing?
Content-based routing = enviar subconjuntos del stream a distintos destinos según el contenido de los datos. La idea: no todo el stream va al mismo sitio; según lo que dice cada evento, se decide su destino.
Derived streams
Un derived stream es una versión transformada del stream original que habilita el content-based routing. Piénsalo como un "sub-stream filtrado" que surge de aplicar transformaciones al stream original.
Escenario: los sensores reportan temperatura cada segundo y quieres enviar las temperaturas altas (>80°) a Activator para una alerta inmediata, guardar las medias horarias en una KQL database para analítica y persistir todas las lecturas en el Lakehouse como backup.
Stream: IoT_Sensors_Original
↓
Transform Filter: temperature > 80
↓
Derived stream "HighTemps" ─→ Activator
Original + Aggregate media horaria
↓
Derived stream "HourlyAvg" ─→ Eventhouse
Original (pass-through)
↓
Derived stream "AllReadings" ─→ Lakehouse
Ventajas
⚡ Eficiente: un solo source, varios outputs sin duplicar la ingesta.
🎯 Consciente del contenido: cada destino recibe exactamente lo que necesita.
🔧 Modular: los cambios en un derived stream no afectan a los otros.
🔄 Reutilizable: los derived streams pueden ser sources de otros eventstreams.
🎯 Cae en el examen. Pregunta típica: "You want to route high-value transactions to a real-time alerting system while sending all transactions to a data warehouse for historical analysis. What do you use?" → Respuesta: derived streams para content-based routing dentro de un eventstream.
🎨
18. Patrones comunes de transformación
Pipeline de calidad de datos
Objetivo: filtrar los eventos inválidos o incompletos antes de almacenarlos.
Ejemplo: añadir total = quantity * price, renombrar cust_id → customer_id, convertir date_str → datetime y hacer join con la master data de clientes.
Agregación y resumen
Objetivo: totales móviles, medias y conteos sobre ventanas temporales.
Source → Filter → Group by [con window] → Destination
SELECT device_id, AVG(temp) AS avg_temp
FROM IoT_Stream
GROUP BY device_id, TumblingWindow(minute, 5)
Estandarización de formato
Objetivo: estructura de datos consistente entre sources antes de combinarlos.
Source A → Manage fields (alinear schema) ─┐
├─→ Union → Destination
Source B → Manage fields (alinear schema) ─┘
Detección de anomalías IoT
Objetivo: filtrar → agregar → enrutar según umbral.
IoT Source
↓
Filter (descartar errores de sensor)
↓
Manage fields (añadir campo de prioridad)
↓
Group by (media horaria por ubicación)
↓
Derived stream (anomalías) → Activator
↓
Datos normales → Lakehouse
🌊
19. DeltaFlow: CDC pipelines simplificados
¿Qué es DeltaFlow?
DeltaFlow (Preview) es una feature de Eventstream que simplifica los pipelines de CDC. En vez de montar un pipeline complejo con Event Hubs + Stream Analytics para procesar los eventos de CDC, DeltaFlow lo hace todo en la configuración del conector del source.
Pipeline CDC tradicional (antes)
Azure SQL DB
↓ (CDC)
Event Hubs
↓
Stream Analytics job:
- Parsear el JSON de Debezium
- Reestructurar
- Gestionar la evolución del schema
↓
Destination
Problemas: complejo, propenso a errores y con mucho código de pegamento.
Con DeltaFlow (después)
Azure SQL DB (o PostgreSQL, MySQL, SQL Server on VM)
↓ [configuración DeltaFlow]
↓ Output analytics-ready
↓ Tablas destino creadas automáticamente
↓ Evolución de schema automática
Eventhouse
Ventajas: ✅ sin código de pegamento, ✅ creación automática de las tablas destino, ✅ gestión automática de la evolución del schema y ✅ configuración sencilla en el conector.
Sources compatibles con DeltaFlow
Azure SQL Database
PostgreSQL
SQL Server on VM
Azure SQL Managed Instance
Feature moderna: Eventstreams SQL para CDC
Nueva feature de febrero de 2026: procesar streams de CDC usando el SQL de Fabric Eventstreams. Ahora puedes transformar los eventos CDC raw de las bases de datos con SQL familiar, centralizando la lógica de shaping del CDC en el eventstream en vez de empujar la complejidad a los consumidores downstream.
SELECT
payload.op AS operation, -- 'c' create, 'u' update, 'd' delete
payload.after.customer_id,
payload.after.email,
payload.after.updated_at
FROM CDC_Source
WHERE payload.op IN ('c', 'u')
Feature de febrero de 2026: la integración entre Spark Structured Streaming, los notebooks y Real-Time Intelligence. Combina lo mejor de dos mundos: Eventstream (sources, routing, dashboards de RTI) y Spark Structured Streaming (código PySpark potente, ML, transformaciones complejas).
Capabilities principales
Descubrir streaming sources en el Real-Time Hub.
Código PySpark autogenerado para consumir esos sources.
Reutilizar notebooks existentes como procesadores.
Conectar de forma segura vía autenticación con Entra ID.
Ejemplo mental
Antes (separado): un Eventstream para la ingesta, un notebook de Spark aparte para el ML y "pegamento" manual entre ambos.
Ahora (integrado): en el Real-Time Hub descubres el streaming source, con un click autogeneras el código PySpark que lee del stream, aplicas ML en el notebook y los resultados van a Eventstream/Eventhouse.
🎯 Cuándo usar Spark Structured Streaming vs Eventstream
Eventstream ✅ cuando las transformaciones se expresan en SQL/no-code, no necesitas ML y el equipo es low-code.
Spark Structured Streaming (con notebooks) ✅ cuando hay transformaciones complejas en PySpark, ML inline con los datos en streaming, código custom que Eventstream no soporta o integración con librerías de Python.
La combinación ✅ cuando la ingesta y el routing son sencillos (Eventstream) pero el procesamiento es complejo (Spark). Lo mejor de ambos mundos.
📊
21. Monitoring y troubleshooting
Workspace Monitoring para Eventstreams
Feature de agosto de 2026: workspace monitoring de Eventstream con control por eventstream. Ahora puedes elegir qué eventstreams emiten datos de rendimiento, errores y salud de nodos. Los datos van a 3 tablas KQL del Eventhouse de workspace monitoring.
Métricas típicas monitorizadas
Tasa de ingesta (eventos/segundo).
Latencia (source → destination).
Errores (parsing, conexión, throttling).
Salud de los nodos (CPU, memoria).
Throughput por transformación.
Real-Time Hub para eventos de capacidad
Capacity Overview Events (GA en agosto de 2026): transmite el resumen de capacidad y las señales de estado para monitorizar la utilización y la salud, detectar throttling o cambios de ciclo de vida, y disparar alertas o acciones automatizadas. Muy útil para no llevarse sorpresas con los límites de capacidad.
Troubleshooting común
Problema: no llegan datos al destino. Comprobaciones:
¿El source está funcionando? → ver la preview del source.
¿Las transformaciones son válidas? → ver la preview de output de cada transformación.
¿El destino está accesible? → ver el estado de la conexión.
¿Hay throttling? → ver las métricas de capacidad.
Problema: schema mismatch en el destino. Comprobaciones: ¿Manage fields aplicó los cambios esperados? (ver preview) y ¿coinciden los tipos? (comparar los schemas de source y destination).
Problema: latencia alta. Comprobaciones: ¿las ventanas son demasiado grandes? (reducir el tamaño), ¿hay muchas transformaciones complejas? (simplificar o dividir) y ¿el source está saturado? (ver la tasa de ingesta).
⚠️
22. Trampas típicas y confusiones frecuentes
Trampa 1: "Enhanced y Standard son intercambiables"
❌ No exactamente. Enhanced es la versión moderna con más features; Standard existe por compatibilidad legacy. Para eventstreams nuevos → Enhanced.
Trampa 2: "Tumbling y Hopping son lo mismo"
❌ Tumbling NO tiene solapamiento. Hopping SÍ. Una tumbling equivale a una hopping donde hop = size.
Trampa 3: "Un evento solo pertenece a una ventana siempre"
❌ Solo en Tumbling. En Hopping y Sliding un evento puede estar en varias ventanas a la vez.
Trampa 4: "La sliding window emite output a intervalos regulares"
❌ La sliding es event-driven: emite solo cuando cambia el contenido (un evento entra/sale). Las tumbling y hopping SÍ son programadas.
Trampa 5: "La ventana máxima es de 24 horas"
❌ Máximo 7 días para todos los tipos de ventana.
Trampa 6: "La session window no necesita parámetros"
❌ Necesita timeout, maximum duration y una partition opcional. La única sin parámetros es la Snapshot.
Trampa 7: "El content-based routing requiere varios eventstreams"
❌ Un solo eventstream con derived streams puede enrutar a varios destinos según el contenido.
Trampa 8: "Puedo usar Fabric Activator sin conectarlo desde el eventstream"
❌ Activator puede ser un destino del eventstream, no un paso intermedio. Los eventos fluyen a Activator y sus rules deciden las acciones.
Trampa 9: "Union combina cualquier campo"
❌ Union solo mantiene los campos compartidos (mismo nombre y tipo). Los que no coinciden se descartan.
Trampa 10: "Un eventstream = un destino"
❌ Varios destinos en paralelo desde el mismo eventstream. Es una feature clave del patrón.
Trampa 11: "El SQL operator soporta cualquier SQL"
❌ Soporta un subconjunto de SQL específico para streaming (el mismo que Stream Analytics). No es T-SQL completo.
Trampa 12: "DeltaFlow sirve para cualquier base de datos"
❌ Solo para Azure SQL Database, PostgreSQL, SQL Server on VM y Azure SQL Managed Instance.
Trampa 13: "La snapshot window agrupa eventos por segundo/minuto"
❌ La snapshot agrupa los eventos con el MISMO timestamp exacto (system time). No por bucket temporal.
Trampa 14: "El output de una ventana contiene todos los eventos"
❌ El output es un único evento por ventana con la agregación calculada. Los eventos individuales no se emiten.
Trampa 15: "Puedo aplicar varios windowing en el mismo GROUP BY"
✅ ¡Sí! Con la función Windows() puedes agregar sobre varias ventanas temporales en el mismo GROUP BY.
🚀 ¡Módulo 9 completado! Ya dominas las windowing functions y el content-based routing, que son de lo más examinable del curso. Sigue con el Módulo 10: Eventhouse y KQL. 🌸