DP-700: Implementar pipelines de ingestión con Dataflows y Synapse
Voy a enseñar cómo implementar pipelines de ingestión de datos en el contexto de Fabric (Dataflows / Synapse pipelines), una competencia central en el DP-700. Saber esto ayuda en el examen y, en la práctica, garantiza que los datos lleguen fiables y de forma eficiente a su entorno de analytics. Aquí encontrará conceptos, pasos concretos, ejemplos con cifras plausibles y buenas prácticas para producción.
Qué necesita saber
Ingestión de datos es el proceso de traer datos desde orígenes (ficheros, bases de datos, servicios) hacia su capa de almacenamiento/procesamiento (OneLake / data lake o tabular storage). En Fabric puede usar Dataflows (Power Query) para transformaciones en el ingest path o Synapse pipelines para orquestación y copia de datos. Los conceptos esenciales son:
- Conectores y autenticación: saber configurar credenciales (Managed Identity, service principal, clave) para orquestar copias seguras. Por ejemplo, para un almacenamiento Azure Blob preferir Managed Identity del workspace; para fuentes on‑premises usar un gateway junto con service principal o credenciales seguras.
- Modos de ingestión: carga completa vs incremental; cuándo usar copy incremental con watermark o CDC (Change Data Capture). Para una tabla de 100M de filas, una carga completa diaria es costosa — un incremental que copie 1–5% de los datos reduce costes y tiempo de procesamiento.
- Transformación ligera vs pesada: usar Dataflows/Power Query para limpiezas y mapeo próximos a la fuente (ideal para ficheros CSV/Excel y reglas de calidad), usar compute (Spark, Synapse SQL Pools, Mapping Data Flows) para transformaciones pesadas en datasets de decenas a cientos de GB o más.
- Idempotencia y reprocesamiento: diseñar pipelines que puedan reejecutarse sin duplicar datos. Técnicas comunes: usar claves naturales para deduplicación, upsert/merge en el sink, o mantener logs de ejecución con watermarks. En producción, establecer un periodo de retención (ej.: 90 días) para permitir reprocesamientos sin acumular costes.
Ejemplo simple: copiar un fichero CSV de 2 GB almacenado en un Azure Blob a una tabla en Fabric usando un pipeline de copia. El pipeline necesita el conector Blob, el mapeo de columnas y la política de fallo/retries (por ejemplo, 3 intentos con backoff exponencial). Una copia paralela bien configurada puede alcanzar 100–200 MB/s, reduciendo el tiempo de ingestión a decenas de segundos/minutos, dependiendo de la red y del número de ficheros.
Cómo funciona — paso a paso práctico
A continuación se muestran pasos prácticos para crear un pipeline de ingestión simple que copia datos desde un storage a una tabla en Fabric y ejecuta transforms ligeros con Dataflow. Incluyo sugerencias operativas y valores típicos que puede ajustar según su entorno.
-
Preparar credenciales y linked services: configure un Linked Service para el Azure Blob/ADLS usando Managed Identity del workspace o un service principal. Esto evita colocar secretos en pipelines. En entornos de producción, establecer rotación de credenciales y auditar accesos. Típicamente, se asigna solo el mínimo privilegio (RBAC) a la identidad, por ejemplo, lectura/escritura solo en el container necesario.
-
Crear el Dataflow (Power Query): en Data Factory / Fabric Dataflows cree un flujo que lea el fichero, detecte el separador, defina tipos y aplique reglas de limpieza (trim, reemplazar nulos, normalizar fechas). Esto reduce errores derivados de esquemas inconsistentes. Ej.: transformar una columna de texto a date con fallback a NULL cuando el formato falla, o aplicar una regla de deduplicación por CustomerID manteniendo la fila con mayor timestamp.
// Exemplo conceptual de transformações Power Query Table.ReplaceValue(Source, null, "", Replacer.ReplaceValue, {"CustomerName"}) Table.TransformColumnTypes(PrevStep, {{"OrderDate", type date}}) -
Crear pipeline de copia: añada una actividad Copy Data que use el Dataflow como source o conecte directamente el fichero como fuente y la tabla de Fabric como sink. Configure mapeo de columnas, paralelismo (Degree of Copy Parallelism) y la política de pre‑copy (ej.: truncate vs append). Para cargas regulares, preferir append + upsert/merge en el sink para evitar ventanas de downtime.
-
Implementar ingestión incremental: cuando sea posible, use watermark o columna de modificación para copiar solo filas nuevas/alteradas. Configure la query source para filtrar por timestamp > @pipelineVariable('lastWatermark') y actualice la variable al final de la ejecución. Por ejemplo, un watermark almacenado en un fichero JSON o tabla de metadatos actualizado al final de cada ejecución. Esto reduce el volumen transferido en cargas grandes de terabytes.
// Padrão conceptual de filtro incremental SELECT * FROM SourceTable WHERE ModifiedAt > @pipeline().parameters.lastWatermark -
Planificación y monitorización: programe el pipeline (trigger por tiempo — ej.: cada 15 minutos, diario — o event trigger cuando llega un fichero). Active retry (por ejemplo, 3 intentos con 30s, 60s, 120s) y notificaciones (email/Teams) en caso de fallo. Monitorice métricas: número de filas, bytes transferidos, duración y tasa de error. Establezca SLAs — ej.: ingestión diaria completada en 2 horas para datasets de 500 GB.
-
Probar idempotencia: ejecute el pipeline varias veces con la misma entrada para garantizar que no duplica datos — use claves naturales o lógica de upsert en el sink (MERGE). En escenarios OLTP, un enfoque común es aplicar MERGE por batch_id o ModifiedAt para garantizar consistencia.
Errores comunes
- Credenciales mal configuradas: usar claves/strings en vez de Managed Identity aumenta el riesgo y provoca fallos cuando los secretos expiran. Siempre preferir identidades gestionadas cuando sea posible y auditar errores de autenticación en el log para detección temprana.
- Descuidado con los esquemas: asumir que todos los ficheros tienen el mismo esquema causa fallos o corrompe datos. Validar esquema y tratar columnas ausentes/extra es esencial. En entornos con cientos de fuentes, crear reglas de validación automática que rechacen ficheros fuera de lo esperado reduce el riesgo operativo.
- Ingestión completa cuando debería ser incremental: ejecutar cargas completas indiscriminadamente aumenta coste y tiempo; definir watermark/CDC es frecuentemente la mejor opción para datos grandes. Por ejemplo, sustituir una carga diaria completa de 1 TB por un incremental de 20 GB reduce el coste de transferencia y el tiempo de procesamiento drásticamente.
Cómo practicar
Practicar con ejercicios en el propio entorno Fabric es fundamental. Cree labs que simulen ficheros de 100 MB a 5 GB, implemente triggers basados en eventos y pruebe reejecuciones y fallos simulados. Use el Practice Assessment OFICIAL de Microsoft (gratuito) para evaluar las áreas donde necesita reforzar conocimiento y consulte la study guide oficial de Microsoft (gratuita) para los temas medidos. No utilice dumps ni preguntas de examen no oficiales — practique con labs y los recursos oficiales.
En resumen
- La ingestión combina conectores, autenticación, modos (completo vs incremental) y transformación ligera/pesada.
- Use Managed Identity y Linked Services para seguridad y fiabilidad de las pipelines.
- Implemente watermark/CDC para reducir coste y tiempo en ingestiones incrementales; considere políticas de retención y ventanas de retención para reprocesamiento.
- Pruebe idempotencia y validación de esquema para evitar duplicación y corrupción de datos; monitorice métricas y defina SLAs claros.