Cómo hacer limpieza incremental de archivos CSV en Azure Data Factory
Este tutorial explica cómo implementar una pipeline en Azure Data Factory para detectar archivos CSV nuevos o modificados, validar/limpiar los datos y cargar solo los registros incrementales en una tabla en Azure SQL. Es útil para reducir costes y evitar duplicación al procesar ingestas recurrentes de archivos.
Requisitos previos
- Cuenta Azure con permisos para crear recursos.
- Un recurso Azure Data Factory v2 creado.
- Almacenamiento Blob o ADLS Gen2 con algunos archivos CSV de ejemplo.
- Una base de datos Azure SQL con tabla destino y acceso de escritura.
- Conocimientos básicos de pipelines, datasets y activities en Azure Data Factory.
Paso 1: Concepto — cómo funciona la ingestión incremental
Explicación sencilla: usamos la lista de archivos (Get Metadata) para detectar archivos nuevos/alterados en base a lastModified. Guardamos un registro del último procesamiento (watermark) y solo procesamos archivos con lastModified posterior. Esto evita reprocesar archivos ya validados.
Paso 2: Crear un archivo de control (watermark) en Blob
Vamos a guardar la fecha/hora del último procesamiento en un pequeño archivo JSON en el blob. Cree un archivo llamado watermark.json con contenido inicial:
{
"lastRun":"1970-01-01T00:00:00Z"
}
Paso 3: Pipeline — obtener lista de archivos nuevos
En ADF cree una pipeline con las siguientes actividades: Get Metadata para listar archivos, Lookup para leer el watermark.json, Filter para seleccionar solo los archivos con lastModified > watermark.
// Get Metadata dataset: apuntar a la carpeta de CSVs, field: childItems
// Lookup dataset: apuntar a watermark.json
// Ejemplo de expresión en el Filter activity para comparar fechas
@greater(item().lastModified, pipeline().parameters.lastWatermark)
Paso 4: Leer el lastRun del watermark (Lookup)
Añada un Lookup que lea watermark.json. En el panel Settings, active First row only. Guarde el valor en una variable de pipeline llamada lastWatermark usando un Set Variable con la expresión:
@activity('LookupWatermark').output.value.lastRun
Paso 5: Usar Get Metadata para obtener lastModified de los archivos
Get Metadata con dataset de la carpeta (Field list: Child Items). La salida da los nombres; para cada item se quiere el lastModified — puede usar un Lookup por archivo o, preferible, habilitar dataset parametrizado y usar una actividad ForEach sobre los childItems para ejecutar un Get Metadata individual que devuelva lastModified.
// Dentro del ForEach (items: @activity('GetFolder').output.childItems)
// Get Metadata (param: fileName) -> field: lastModified
// Exponer output: activity('GetFileMetadata').output.lastModified
Paso 6: Filtrar archivos nuevos y limpiar datos simples
Dentro del ForEach, tras obtener lastModified, use una actividad If Condition para probar si lastModified > lastWatermark. Si true, ejecutar un Data Flow o Copy con mapping y validaciones simples (ej.: eliminar filas con campos obligatorios vacíos, normalizar fechas).
// If Condition expression
@greater(formatDateTime(activity('GetFileMetadata').output.lastModified,'yyyy-MM-ddTHH:mm:ssZ'), variables('lastWatermark'))
// Ejemplo de transformaciones simples en un Mapping Data Flow:
// - Source: CSV
// - Derived Column: trim() y parseDate()
// - Filter: isNotNull(campo_clave)
// - Sink: Azure SQL (modo upsert con clave única)
Paso 7: Carga incremental a Azure SQL
Para evitar duplicados, use en el Sink del Data Flow la opción de update/insert (upsert) basada en una columna clave (por ejemplo, id o combinación de campos). Alternativa: cargar a una tabla staging y ejecutar stored procedure para deduplicar con MERGE.
// Ejemplo minimal de MERGE (Azure SQL) para deduplicar tras carga en staging
MERGE dbo.Target AS T
USING dbo.Staging AS S
ON T.Key = S.Key
WHEN MATCHED THEN UPDATE SET T.Col = S.Col
WHEN NOT MATCHED THEN INSERT (Key, Col) VALUES (S.Key, S.Col);
Paso 8: Actualizar el watermark tras el éxito
Al final de la pipeline, tras la confirmación de éxito de la carga, escriba en el watermark.json la fecha/hora máxima de los archivos procesados (por ejemplo, now()). Use una actividad Web para llamar a la REST API del Blob (o una actividad Copy con dataset JSON de salida) para sobrescribir el archivo.
// Ejemplo expresión para nuevo watermark
@utcNow() // o máximo entre archivos procesados
// Para escribir con Copy activity: fuente = una pequeña tabla/variable, sink = dataset watermark.json
Verificar el resultado
Valide: 1) la pipeline se ejecutó sin errores; 2) la tabla destino contiene solo los registros esperados; 3) watermark.json fue actualizado con la nueva fecha; 4) re-ejecutar la pipeline no reprocesa archivos ya procesados. Use Monitor de Azure Data Factory y consultas en Azure SQL para ver cambios.
Conclusión
Con este patrón consigue procesar archivos CSV de forma incremental, validar y cargar en Azure SQL reduciendo coste y duplicación. Próximos pasos: añadir logging detallado, gestionar errores con Dead-letter (staging) y parametrizar para varias carpetas. Consejo: empiece por probar con pocos archivos y verifique siempre el timezone de las fechas.