Prueba Técnica – SQL (Stored Procedures) y Databricks
Solución propuesta
1. Procedimiento almacenado [dbo].
[BEESP_SYNC_SnowFlake_Tracks]
A continuación se listan los errores encontrados en el SP, agrupados en errores de sintaxis (impiden la
compilación/ejecución) y errores de lógica de negocio (compilan pero producen resultados incorrectos), seguidos de las
buenas prácticas recomendadas para su implementación.
1.1 Errores de sintaxis (el procedimiento no compila)
a) OPENQUERY del bloque ANDROID sin filtro de fecha
Código original:
WHERE CAST(SENT_AT AS DATE) >=
GROUP BY
Corrección propuesta:
WHERE CAST(SENT_AT AS DATE) >= CAST(DATEADD(DAY,-1,GETDATE()) AS DATE)
GROUP BY
Al comparado del bloque IOS (línea equivalente sí tiene el valor), en el bloque ANDROID el operador ">=" queda sin el lado
derecho de la expresión. Esto rompe la cadena de texto que se envía por OPENQUERY al linked server de Snowflake y el
procedimiento no compila.
b) CAST sin tipo de dato (ISPID)
Código original:
CAST(NULL AS ) AS ISPID,
Corrección propuesta:
CAST(NULL AS INT) AS ISPID,
CAST exige un tipo de dato explícito; dejarlo vacío es un error de sintaxis. Además, más adelante el UPDATE sí asigna
[Link] = [Link], por lo que la columna debe existir con un tipo consistente con VPV_WEB.dbo.BEES_ISP.ISPID
(típicamente INT).
c) INSERT INTO BEES_Device: coma colgante y columna faltante en la lista de destino
Código original:
INSERT INTO VPV_WEB.dbo.BEES_Device
(
DeviceName,
DeviceModel
)
SELECT ISNULL(TMP.CONTEXT_DEVICE_MANUFACTURER,''), ISNULL(TMP.CONTEXT_DEVICE_NAME,''),
ISNULL(TMP.CONTEXT_DEVICE_MODEL,''),
FROM #BEES_SnowFlake_Tracks AS TMP
Corrección propuesta:
INSERT INTO VPV_WEB.dbo.BEES_Device
(
DeviceBrand,
DeviceName,
DeviceModel
)
SELECT ISNULL(TMP.CONTEXT_DEVICE_MANUFACTURER,''), ISNULL(TMP.CONTEXT_DEVICE_NAME,''),
ISNULL(TMP.CONTEXT_DEVICE_MODEL,'')
FROM #BEES_SnowFlake_Tracks AS TMP
La lista de columnas del INSERT solo trae DeviceName y DeviceModel, pero el SELECT devuelve 3 columnas
(Manufacturer, Name, Model) y además termina con una coma colgante antes de FROM. La cláusula WHERE NOT EXISTS
más abajo sí compara [Link], lo que confirma que falta esa columna en la lista de inserción. Se agrega
DeviceBrand y se quita la coma final.
d) ISNUMERIC sin operador de comparación
Código original:
WHERE ISNUMERIC(TMP.CONTEXT_TRAITS_POC_ID) 1
Corrección propuesta:
WHERE TRY_CAST(TMP.CONTEXT_TRAITS_POC_ID AS INT) IS NOT NULL
Falta el operador "=" entre la función y el 1 ("ISNUMERIC(...) 1" no es una expresión válida). Adicionalmente,
ISNUMERIC no es confiable para validar enteros (acepta '.', '$', '1e3', etc.), por lo que se recomienda reemplazarlo por
TRY_CAST/TRY_CONVERT, que además evita el error de conversión descrito en el punto (h).
e) NOT EXISTS con predicado incompleto (expCodigo)
Código original:
WHERE [Link] = @TimePeriodID
AND [Link] = )
Corrección propuesta:
WHERE [Link] = @TimePeriodID
AND [Link] =
[Link])
La condición "[Link] = )" no tiene el lado derecho de la comparación; el paréntesis de cierre queda huérfano. Debe
correlacionarse con la tabla temporal (por ejemplo [Link], o el campo equivalente de código de cliente) para que
el DELETE elimine únicamente los registros cuyo cliente no exista en STR_Expendio.
f) UPDATE sin la palabra clave SET (dos casos)
Código original:
UPDATE Bee
[Link] = [Link]
FROM #BEES_SnowFlake_Tracks AS Bee
...
UPDATE Bee
[Link] = [Link]
FROM #BEES_SnowFlake_Tracks AS Bee
Corrección propuesta:
UPDATE Bee
SET [Link] = [Link]
FROM #BEES_SnowFlake_Tracks AS Bee
...
UPDATE Bee
SET [Link] = [Link]
FROM #BEES_SnowFlake_Tracks AS Bee
Los dos primeros UPDATE del bloque de actualización de llaves foráneas no incluyen SET (solo el tercero, el de ISPID, lo
tiene). Sin SET la asignación de columna es sintácticamente inválida.
g) INSERT INTO BEES_Tracks: columna vacía y coma final
Código original:
INSERT INTO BEES_Tracks
(
Date, CustomerID, IpAddress, Context_DeviceID, DeviceID, PlatFormID,
,
Phone, Q_Events, Fraud,
)
Corrección propuesta:
INSERT INTO BEES_Tracks
(
Date, CustomerID, IpAddress, Context_DeviceID, DeviceID, PlatFormID,
Email,
Phone, Q_Events, Fraud
)
Hay una posición vacía en la lista de columnas (debe ser "Email", ya que el SELECT trae MAX(Context_Traits_Email) AS
Email) y una coma sobrante después de Fraud, justo antes del paréntesis de cierre. Ambas son errores de sintaxis.
1.2 Errores de lógica / diseño (el SP podría compilar, pero el resultado es incorrecto o
frágil)
h) CAST a INT antes de filtrar por ISNUMERIC — riesgo de error en tiempo de ejecución
Código original:
CAST(TMP.CONTEXT_TRAITS_POC_ID AS INT) AS CustomerID,
...
WHERE ISNUMERIC(TMP.CONTEXT_TRAITS_POC_ID) 1
Corrección propuesta:
TRY_CAST(TMP.CONTEXT_TRAITS_POC_ID AS INT) AS CustomerID,
...
WHERE TRY_CAST(TMP.CONTEXT_TRAITS_POC_ID AS INT) IS NOT NULL
SQL Server no garantiza evaluar el WHERE antes que el SELECT (el optimizador puede reordenar). Aunque el filtro elimine
los no numéricos, el CAST en el SELECT puede fallar antes con "Conversion failed when converting the varchar value ... to
int". La forma segura es usar TRY_CAST en ambos lugares.
i) ISDATE() usado como si convirtiera la fecha
Código original:
DATEADD(HOUR, -5, ISDATE(LEFT(SENT_AT,19))) AS SENT_AT,
Corrección propuesta:
DATEADD(HOUR, -5, CAST(LEFT(SENT_AT,19) AS DATETIME)) AS SENT_AT,
ISDATE() solo valida si una cadena es una fecha válida y devuelve 0/1 (BIT); no la convierte a datetime. El resultado real de
esta expresión es siempre "DATEADD(HOUR,-5, 1)" (o 0), es decir, un entero interpretado como fecha base de SQL Server,
no la fecha del evento. Debe usarse CAST/CONVERT (o mejor, TRY_CONVERT) sobre el texto.
j) Bandera de fraude que nunca se "apaga"
Código original:
UPDATE Btk
SET
[Link] = 1
FROM BEES_Tracks Btk INNER JOIN #Fraud Frd ON Btk.Context_DeviceID = Frd.Context_DeviceID
Corrección propuesta:
UPDATE Btk
SET [Link] = CASE WHEN Frd.Context_DeviceID IS NOT NULL THEN 1 ELSE 0 END
FROM BEES_Tracks Btk LEFT JOIN #Fraud Frd ON Btk.Context_DeviceID = Frd.Context_DeviceID
El proceso solo marca Fraud = 1 para los dispositivos que cumplen la condición actual (más de 3 clientes distintos en 21
días), pero nunca vuelve a poner Fraud = 0 si el dispositivo deja de cumplir la condición en una ejecución posterior. Con un
LEFT JOIN y un CASE se recalcula el estado en cada corrida.
1.3 Buenas prácticas recomendadas
• Agregar SET NOCOUNT ON al inicio del SP para evitar mensajes de conteo de filas que afectan el desempeño y
pueden interferir con clientes/ETL que consumen el SP.
• Envolver toda la lógica en TRY…CATCH con manejo de transacción (BEGIN TRAN / COMMIT / ROLLBACK).
Actualmente, si falla el TRUNCATE, un OPENQUERY o cualquier INSERT/UPDATE intermedio, el proceso queda a
medias sin registro del error ni posibilidad de revertir cambios ya aplicados en las tablas destino (BEES_Platform,
BEES_Device, BEES_ISP, BEES_Tracks).
• Evitar TRUNCATE + reconstrucción total de la tabla de staging combinado con una ventana fija de 1 día
(DATEADD(DAY,-1,GETDATE())). Si el job no corre un día (feriado, falla, mantenimiento), esos datos se pierden
para siempre. Es preferible una estrategia incremental basada en "marca de agua" (última fecha cargada) o reprocesar
una ventana de solape (p. ej. últimos 3 días) con upsert/MERGE idempotente.
• Eliminar la duplicación de código entre los bloques ANDROID / IOS / WEB (prácticamente idénticos). Se puede
resolver con una tabla de configuración de orígenes (nombre de fuente, base Snowflake, plataforma) recorrida en un
bucle, o unificando la extracción con UNION ALL parametrizado, reduciendo el riesgo de que una corrección se
aplique en un bloque y se olvide en otro (como pasó con el filtro de fecha faltante en ANDROID).
• Especificar siempre la lista explícita de columnas en los INSERT INTO ... SELECT (los dos primeros bloques insertan
sin lista de columnas), para no depender del orden físico de columnas de la tabla destino.
• Sustituir el patrón manual INSERT ... WHERE NOT EXISTS para las tablas de dimensión (BEES_Platform,
BEES_Device, BEES_ISP) por una sentencia MERGE, más clara, atómica y menos propensa a condiciones de carrera
si el SP llegara a ejecutarse en paralelo.
• Revisar el uso de NOLOCK en los JOIN de dimensiones: permite lecturas sucias (dirty reads) y puede traer filas
duplicadas/perdidas si hay inserciones concurrentes; si el volumen lo permite, preferir READ COMMITTED
SNAPSHOT a nivel de base de datos en lugar de NOLOCK disperso en el código.
• Uniformar el tipo de dato usado para columnas de "llave de negocio" (por ejemplo Context_Device_ID se define como
VARCHAR(80) al leer de Snowflake pero luego se castea a CHAR(20) en la tabla temporal); CHAR rellena con
espacios y puede romper comparaciones de igualdad en los JOIN y UPDATE posteriores.
• Crear índices sobre la tabla temporal #BEES_SnowFlake_Tracks en las columnas usadas para JOIN/UPDATE
(CONTEXT_OS_NAME, CONTEXT_DEVICE_MANUFACTURER/NAME/MODEL,
CONTEXT_TRAITS_POC_ID) ya que hoy solo se indexa #Fraud.
• Parametrizar los valores fijos del proceso de fraude (ventana de 21 días, umbral de 3 clientes) como variables o
parámetros del SP en lugar de dejarlos como literales embebidos.
• Agregar registro/log de auditoría (filas insertadas/actualizadas, tiempo de ejecución, errores) similar al patrón de
MonitoringService que sí se usa en el notebook de PySpark, para tener trazabilidad de estas cargas también en SQL
Server.
2. Azure Data Factory
¿Qué es Azure Data Factory y cuál es su función principal?
Es el servicio de integración de datos de Azure para orquestar y automatizar el movimiento y la transformación de datos
(ETL/ELT) entre orígenes on-premises y en la nube, sin necesidad de escribir infraestructura de servidores. Su función
principal es construir, programar y monitorear pipelines (canalizaciones) que conectan, mueven, transforman y publican
datos hacia destinos como Data Lake, Synapse, Databricks, SQL Server, entre otros.
¿Cómo se configuran las actividades de copia en Data Factory?
• Se crean Linked Services que definen la conexión/credenciales al origen y al destino (base de datos, storage account,
API, etc.).
• Sobre cada Linked Service se definen Datasets, que describen el esquema/estructura de los datos a leer o escribir (tabla,
contenedor, ruta de archivo, formato).
• Dentro de un pipeline se agrega una actividad Copy Data, indicando el dataset de origen (source) y el dataset de destino
(sink), el mapeo de columnas (schema mapping), y opciones de rendimiento como paralelismo (degree of copy
parallelism), tamaño de lote (batch size) y política de reintentos/tolerancia a fallas.
• Opcionalmente se agregan filtros/consultas (query, filtro por partición o marca de tiempo) en el origen para soportar
cargas incrementales.
¿Cuál es la diferencia entre un flujo de datos y un conjunto de datos en Data Factory?
Un Dataset (conjunto de datos) es solo la representación de la estructura de los datos en un almacén (por ejemplo, una tabla o
un archivo) apuntando a un Linked Service; no ejecuta lógica, solo describe "qué" y "dónde" están los datos. Un Data Flow
(flujo de datos) es un componente de transformación visual (basado en Spark administrado por ADF) donde se diseñan pasos
de transformación —joins, agregaciones, pivots, limpieza, validaciones— sobre uno o varios datasets de entrada, generando
datasets de salida. En otras palabras, el dataset es la referencia a los datos y el data flow es la lógica de transformación que se
aplica sobre ellos dentro de un pipeline.
¿Cómo se monitorean y depuran las canalizaciones en Data Factory?
• Desde la pestaña Monitor de ADF Studio se visualizan las ejecuciones de pipelines y triggers, con estado (Succeeded,
Failed, In Progress), duración, y la vista Gantt/consumo de cada actividad.
• Se puede entrar al detalle de cada actividad para ver inputs/outputs, mensajes de error y código de error específico del
conector.
• Para depurar antes de publicar, se usa el modo Debug del pipeline (ejecución de prueba sin necesidad de trigger) y, en
Data Flows, el Data Flow Debug con vista de datos y estadísticas de cada transformación.
• Se pueden configurar alertas (Azure Monitor / Log Analytics) sobre fallas de pipeline, y activar reintentos automáticos
y políticas de tiempo de espera por actividad.
3. Databricks (Clusters)
¿Qué es Databricks y cómo se relaciona con Apache Spark?
Databricks es una plataforma unificada de datos y analítica construida por los creadores de Apache Spark, que ofrece un
entorno administrado (notebooks colaborativos, clusters administrados, orquestación de jobs, catálogo de datos vía Unity
Catalog, MLflow, etc.) sobre un motor de cómputo distribuido basado en Spark optimizado (Databricks Runtime/Photon). Es
decir, Spark es el motor de procesamiento distribuido subyacente, y Databricks es la capa de plataforma que facilita su uso en
producción: administración de clusters, autoescalado, seguridad, almacenamiento (Delta Lake) e integración con el resto del
ecosistema de datos.
¿Cuáles son las ventajas de usar clústeres de Databricks frente a clústeres de Spark
independientes?
• Administración simplificada: creación, autoescalado y terminación automática de clústeres sin gestionar infraestructura
manualmente (a diferencia de un clúster Spark self-managed).
• Optimizaciones propietarias: motor Photon y mejoras del Databricks Runtime que aceleran consultas SQL y cargas de
trabajo de Spark frente a Apache Spark open source puro.
• Delta Lake nativo: transacciones ACID, time travel, esquema evolutivo, comandos MERGE/UPDATE/DELETE
eficientes sobre el data lake.
• Colaboración: notebooks compartidos, control de versiones, comentarios y ejecución multi-lenguaje (SQL, Python,
Scala, R) en el mismo notebook.
• Seguridad y gobierno centralizados con Unity Catalog (permisos a nivel de catálogo/esquema/tabla/columna, linaje de
datos), algo que un clúster Spark independiente no ofrece de fábrica.
• Integración con orquestación (Jobs/Workflows), MLflow para ciclo de vida de modelos, y conectividad simplificada a
fuentes como Snowflake, Kafka, ADLS, S3, etc.
Tres lenguajes que soporta un notebook en Databricks
• Python
• SQL
• Scala (también R, y Shell mediante %sh)
4. Revisión del notebook PySpark (maz_dm_combos)
El notebook incluido en el documento (embebido como objeto OLE, nombre interno maz_dm_combos) contiene 5 celdas:
inicialización de monitoreo (0.5), imports (1.0), cálculo de fechas (2.0), lectura de la tabla de preventa (2.5) y la carga
principal de la tabla dm_combo vía [Link] con múltiples CTE (3.0).
4.1 Aciertos
• Manejo de fechas dinámico y parametrizado: fecha_inicio/fecha_fin se calculan con relativedelta a partir de la fecha
actual, evitando fechas quemadas en el código de negocio (a diferencia del SP en SQL Server).
• Monitoreo desacoplado y tolerante a fallos: la importación y uso de MonitoringService está envuelta en try/except
(monitor_enabled) y se usa una función auxiliar safe_monitor_call, de modo que si el servicio de monitoreo no está
disponible, el notebook igual continúa ejecutando la carga de datos en lugar de detenerse por completo.
• Uso de [Link] para publicar execution_uuid y proceso_id, permitiendo que otras tareas del mismo
Job (orquestación de Workflows) consuman esos valores; buena práctica de observabilidad end-to-end.
• La consulta principal usa CTEs (usuarios, combos, combinacion) que hacen legible una transformación compleja con
múltiples joins, en vez de una única consulta monolítica.
• Uso de QUALIFY ROW_NUMBER() OVER (...) para deduplicar usuarios por usuario_id, más eficiente y legible que
un subquery con filtro adicional.
• La carga final usa INSERT INTO ... REPLACE WHERE (semántica de Delta Lake), que sobrescribe selectivamente
solo la partición/rango afectado en vez de reescribir toda la tabla, haciendo la carga incremental e idempotente (se puede
re-ejecutar sin duplicar datos).
• El bloque principal está envuelto en try/except con notificación de fallas (monitor.end_process(status='FAILED') y
send_notification), dando trazabilidad básica de errores del proceso.
4.2 Mejoras posibles
• Bug en la celda 2.5: la línea preventa = [Link](...).filter(fecha_inicio <= col('periodo_fin')) & (col('periodo_fin') <=
fecha_fin) cierra el paréntesis de .filter() antes de aplicar el AND (&); el resultado no es un DataFrame filtrado
correctamente sino una operación inválida entre un DataFrame y una Column. Debe ser .filter((col('periodo_fin') >=
lit(fecha_inicio)) & (col('periodo_fin') <= lit(fecha_fin))), envolviendo ambas condiciones dentro del filter y usando lit()
para los valores Python.
• Bug en la misma celda: [Link]("overwrite").saveAsTable(...) referencia una variable df que nunca fue definida
(la variable creada es preventa); tal como está, la celda lanzaría NameError en tiempo de ejecución.
• El mode("overwrite") reescribe toda la tabla tmp_dm_preventa en cada corrida; si la tabla es grande conviene usar
escritura por partición (partitionOverwriteMode=dynamic o replaceWhere), igual que se hizo correctamente en la celda
3.0 para dm_combo.
• Imports no utilizados: numpy (np) y pickle (pk) se importan en la celda 1.0 pero no se usan en ninguna celda visible;
conviene limpiarlos para reducir dependencias y tiempo de arranque del cluster.
• fecha_fin se calcula pero no se usa como cota superior en la consulta SQL principal (solo se filtra
[Link]::date >= '{fecha_inicio}'::date, sin límite superior). Si la intención era procesar una ventana
cerrada [fecha_inicio, fecha_fin], falta agregar la condición superior; si es intencional (ventana abierta hacia el
presente), conviene dejar un comentario explicando por qué fecha_fin no se usa, para evitar que alguien lo interprete
como un olvido en el futuro.
• Interpolación de fechas directamente en el string SQL vía f-string ('{fecha_inicio}'::date). Aunque el riesgo de
inyección es bajo porque los valores vienen de [Link]() y no de un input externo, es más robusto y explícito usar
parámetros con [Link](query, args={...}) (parámetros con nombre, soportado en Databricks Runtime reciente) en vez
de interpolación de texto.
• En el except final, monitor.send_notification(execution_id) se llama directamente (no a través de safe_monitor_call). Si
el error ocurrió precisamente porque monitor nunca se inicializó (monitor_enabled = False), esta línea lanzaría un nuevo
error (NameError/AttributeError) que enmascara la excepción original antes de llegar al raise e. Debe canalizarse
también por safe_monitor_call('send_notification', execution_id).
• La celda de monitoreo mezcla el magic %cd dentro de un bloque try de Python indentado; los magics de celda en
Databricks no siempre se comportan de forma consistente cuando están anidados dentro de estructuras de control. Es
más seguro usar [Link]('/Workspace/Users/.../utils/monitoring') en Python puro.
• La ruta del módulo de monitoreo está codificada al workspace personal de un usuario
(gen_maz_pe_win053@[Link]). Si esa cuenta se deshabilita o cambia de carpeta, todos los notebooks que
dependan de esa ruta fallarán; conviene empaquetar MonitoringService como una librería instalada en el cluster
(wheel/Repos compartido) en lugar de una ruta de usuario.
• La consulta principal hace múltiples LEFT JOIN (cmb_it, cmb_k, dm_material, revenue_maestro_sku, dm_cliente,
pe_portfolio_material, usuarios) más una ventana (SUM ... OVER) y un QUALIFY; conviene revisar el plan de
ejecución y, si las tablas de dimensión son pequeñas, usar hints de broadcast para evitar shuffles costosos a medida que
crezca el volumen.
• El umbral de 10 días de historia (relativedelta(days=-10)) y el desfase horario de -5 horas están hardcodeados; sería
mejor exponerlos como widgets/parámetros del notebook o del Job para poder ajustarlos sin editar código.
4.3 Importancia de los ajustes
• Correctitud: los bugs de la celda 2.5 (variable df inexistente y el filtro mal paréntesis) harían fallar esa celda en
cualquier ejecución real; si el Job no está configurado para detenerse ahí, el notebook seguiría corriendo con datos de
preventa desactualizados o inexistentes en tmp_dm_preventa, afectando a cualquier proceso downstream que dependa
de esa tabla.
• Confiabilidad/observabilidad: si el manejo de errores del monitoreo falla justo cuando ocurre el error real (como en el
punto de send_notification), se pierde la alerta y el equipo se entera del fallo tarde (por reclamos de negocio) en vez de
por el sistema de monitoreo, aumentando el tiempo de detección y resolución de incidentes.
• Rendimiento y costo: reescribir tablas completas con overwrite o no usar broadcast en joins con tablas de dimensión
pequeñas incrementa innecesariamente el tiempo de cluster y, por lo tanto, el costo en Databricks (facturación por
DBU), especialmente si el proceso corre a diario.
• Mantenibilidad: eliminar imports muertos, evitar rutas de usuario hardcodeadas y parametrizar valores (ventana de días,
huso horario) reduce el riesgo de que el proceso se rompa por causas externas (una cuenta deshabilitada, un cambio de
carpeta) y facilita que otro desarrollador entienda y modifique el notebook sin necesidad de contexto tribal.
• Consistencia de datos: aclarar si fecha_fin debe acotar la consulta principal evita ambigüedad sobre si el proceso
realmente procesa una ventana cerrada o simplemente todo lo modificado desde fecha_inicio en adelante, lo cual es
relevante para poder reprocesar rangos históricos de forma predecible.