Nos tomará tres semanas enviar un único canal de datos. Hoy en día, un analista sin experiencia en Python lo hace en un día. Así es como llegamos allí.
Soy Kiril Kazlou, ingeniero de datos de Mindbox. Nuestro equipo recalcula periódicamente las métricas comerciales de los clientes, lo que significa que constantemente creamos data marts para facturación y análisis, extrayendo de docenas de fuentes diferentes.
Durante mucho tiempo, confiamos en PySpark para todo nuestro procesamiento de datos. ¿El problema? Realmente no puedes trabajar con PySpark sin experiencia en Python. Cada nuevo canal requería un desarrollador. Y eso significó esperar, a veces durante semanas.
En esta publicación, le explicaré cómo creamos una plataforma de datos interna donde un analista o gerente de producto puede poner en marcha un proceso actualizado periódicamente escribiendo solo cuatro archivos YAML.
Por qué PySpark nos estaba frenando
Permítanme ilustrar el problema con un ejemplo de libro de texto: calcular MAU (usuarios activos mensuales).
A primera vista, esto parece un trabajo SQL simple: COUNT(DISTINCT customerId) en algunas tablas durante una ventana de tiempo. Pero debido a todos los gastos generales de infraestructura (PySpark, configuración de Airflow DAG, asignación de recursos de Spark, pruebas), tuvimos que entregárselo a los desarrolladores. ¿El resultado? Una semana completa sólo para enviar un contador MAU.
Cada nueva métrica tardó entre una y tres semanas en entregarse. Y cada vez, el proceso fue el mismo:
Un analista definió los requisitos del negocio, encontró un desarrollador disponible y le entregó el contexto. El desarrollador aclaró los detalles, escribió el código PySpark, revisó el código, configuró el DAG y lo implementó.
Lo que realmente queríamos era que los analistas y gerentes de producto (las personas que mejor entienden la lógica empresarial y dominan SQL y YAML) se encargaran de esto ellos mismos. Sin pitón. Sin PySpark.
Con qué reemplazamos PySpark: YAML y SQL son todo lo que necesita
Para adoptar un enfoque declarativo, dividimos nuestra capa de datos en tres partes y elegimos la herramienta adecuada para cada una:
dlt (herramienta de carga de datos): ingiere datos de bases de datos y API externas en el almacenamiento de objetos. Configurado íntegramente a través de un archivo YAML. No se requiere código. dbt (herramienta de creación de datos) en Trino: transforma datos utilizando SQL puro. Vincula modelos a través de ref(), crea automáticamente un gráfico de dependencia y maneja actualizaciones incrementales. Airflow + Cosmos: organiza las tuberías. El DAG Airflow se genera automáticamente a partir de dag.yaml y el proyecto dbt.
Ya estábamos usando Trino como motor de consultas para consultas ad hoc y lo teníamos conectado a Superset para BI. Ya había demostrado su eficacia: para consultas con lógica estándar, procesaba conjuntos de datos masivos más rápido y con menos recursos que Spark. Además de eso, Trino admite de forma nativa el acceso federado a múltiples almacenes de datos desde una única consulta SQL. Para el 90% de nuestros ductos, Trino encajaba perfectamente.
Cómo cargamos datos: dlt.yaml
El primer archivo YAML describe dónde y cómo cargar datos para el procesamiento posterior. A continuación se muestra un ejemplo del mundo real: cargar datos de facturación desde una API interna:
producto: sg-team característica: esquema de facturación: facturación_tarificación dag: dag_id: dlt_billing_tarificación programación: "0 4 * * *" descripción: "Actualización diaria de datos de tarificación" etiquetas: – alertas de facturación: habilitado: verdadero gravedad: fuente de advertencia: tipo: rest_api cliente: base_url: "https://internal-api.example.com" autenticación: tipo: token de portador: dlt-billing.token recursos: – nombre: punto final de tarification_data: ruta: /tarificationData método: POST json: firstPeriod: "{{ fecha_mes_anterior }}" lastPeriod: "{{fecha_mes_anterior }}" pricingPlanLine: CurrentPlan write_disposition: reemplazar pasos_de_procesamiento: – mapa: dlt_custom.billing_tarification_data.map – nombre: cargos_raw columnas: staffUserName: tipo_datos: texto anulable: verdadero punto final: ruta: /data-feed/método de cargos: POST json: primerPeriodo: "{{ fecha_mes_anterior }}" últimoPeriodo: "{{ fecha_mes_anterior }}" escritura_disposición: reemplazar – nombre: descuentos_raw punto final: ruta: /data-feed/método de descuentos: POST json: primerPeriodo: "{{ fecha_mes_anterior }}" últimoPeriodo: "{{ fecha_mes_anterior }}" escritura_disposición: reemplazar
Esta configuración define cuatro recursos de una única API. Para cada uno, especificamos el punto final, los parámetros de solicitud y una estrategia de escritura; en nuestro caso, reemplazar significa "sobrescribir siempre". También puede agregar pasos de procesamiento, definir tipos de columnas y configurar alertas.
La configuración completa tiene 40 líneas de YAML. Sin dlt, cada conector sería un script de Python que maneja solicitudes, paginación, reintentos, serialización al formato de tabla Delta y cargas al almacenamiento.
Cómo transformamos datos con SQL: dbt_project.yaml y sources.yaml
El siguiente paso es configurar el modelo dbt. Con Trino, eso significa consultas SQL.
A continuación se muestra un ejemplo de cómo configuramos el cálculo de MAU. Así es como se ve la preparación de eventos desde una sola fuente:
— int_mau_events_visits.sql (simplificado) {{ config(materialized='table') }} CON período AS ( — Ventana móvil: últimos 5 meses hasta la actualidad SELECCIONE AÑO(CURRENT_DATE – INTERVAL '5' MES) COMO inicio_año, MES(CURRENT_DATE – INTERVAL '5' MES) COMO inicio_mes, AÑO(CURRENT_DATE) COMO fin_año, MES(CURRENT_DATE) AS end_month ), eventos AS ( — Extraiga eventos de visita dentro de la ventana del período SELECT src._tenant, src.unmergedCustomerId, 'visits' AS src_type, src.endpoint FROM {{ source('final', 'customerstracking_visits') }} src CROSS JOIN period p WHERE src.unmergedCustomerId IS NOT NULL AND /* …filtrado de marca de tiempo por límites de año/mes… */ ), events_with_customer AS ( — Resolver ID de clientes fusionados SELECT e._tenant, COALESCE(mc.mergedCustomerId, e.unmergedCustomerId) AS customerId, e.src_type, e.endpoint FROM events e LEFT JOIN {{ ref('int_merged_customers') }} mc ON e._tenant = mc._tenant AND e.unmergedCustomerId = mc.unmergedCustomerId ) – Mantener solo los clientes reales (no eliminados) SELECCIONE ewc._tenant, ewc.customerId, ewc.src_type, ewc.endpoint FROM events_with_customer ewc DONDE EXISTE (SELECCIONE 1 DE {{ ref('int_actual_customers') }} ac DONDE ewc._tenant = ac._tenant Y ewc.customerId = ac.customerId )
Las 10 fuentes de eventos siguen exactamente el mismo patrón. Las únicas diferencias son la tabla fuente y los filtros. Luego los modelos se fusionan en una sola corriente:
— int_mau_events.sql (unión de todas las fuentes) SELECT * FROM {{ ref('int_mau_events_inapps_targetings') }} UNION ALL SELECT * FROM {{ ref('int_mau_events_inapps_clicks') }} UNION ALL SELECT * FROM {{ ref('int_mau_events_visits') }} UNION ALL SELECT * FROM {{ ref('int_mau_events_orders') }} — …más 6 fuentes más
Y finalmente, el data mart donde se agrega todo:
— mau_period_datamart.sql {{ config( materialized='incremental', incremental_strategy='merge', Unique_key=['_tenant', 'start_year', 'start_month', 'end_year', 'end_month'] ) }} {%- set meses_back = var('months_back', 5) | int -%} CON período AS ( SELECCIONAR AÑO (FECHA_CURRENTE – INTERVALO '{{ meses_atrás }}' MES) COMO año_inicio, MES (FECHA_ACTUAL – INTERVALO '{{ meses_atrás }}' MES) COMO mes_inicio, AÑO (FECHA_ACTUAL) COMO año_final, MES (FECHA_ACTUAL) COMO mes_final), eventos_resolvedos AS ( SELECCIONAR * DESDE {{ ref('int_mau_events') }} ), metrics_by_tenant AS ( SELECCIONE er._tenant, COUNT(CASO DISTINTO CUANDO src_type = 'visitas' ENTONCES clientId END) AS CustomersTracking_Visits, COUNT(DISTINCT CASO CUANDO src_type = 'pedidos' ENTONCES customerId END) AS ProcessingOrders_Orders, COUNT(CASE DISTINCT WHEN src_type = 'mailings' THEN customerId END) AS Mailings_MessageStatuses, — …otras métricas COUNT(DISTINCT customerId) AS MAU FROM events_resolved o GROUP BY er._tenant ) SELECT m.*, p.start_year, p.start_month, p.end_year, p.end_month FROM metrics_by_tenant m Periodo de UNIÓN CRUZADA p
Para la configuración del centro de datos, utilizamos incremental_strategy='merge'. dbt genera automáticamente la consulta de combinación, sustituyendo la clave_única por upsert. No es necesario implementar manualmente la carga incremental.
Para vincular los modelos en un solo proyecto, configuramos dbt_project.yaml:
nombre: mau_period versión: '1.0.0' modelos: mau_period: +on_table_exists: reemplazar +on_schema_change: append_new_columns
Y fuentes.yaml, que describe las tablas de entrada:
fuentes: – nombre: base de datos final: data_platform esquema: tablas finales: – nombre: inapps_targetings_v2 – nombre: inapps_clicks_v2 – nombre: customertracking_visits – nombre: Processingorders_orders – nombre: cdp_mergedcustomers_v2 #…
El resultado es la misma lógica de negocios que teníamos en PySpark, pero en SQL puro: sources.yaml reemplaza los esquemas typedspark, {{ ref() }} y {{ source() }} reemplazan .get_table(), y el orden de ejecución automático a través del gráfico de dependencia reemplaza el ajuste manual de recursos de Spark.
Cómo configuramos el flujo de aire: dag.yaml
El cuarto archivo de configuración define cuándo y cómo Airflow ejecuta la canalización:
producto: sg-team característica: esquema de facturación: mau horario: "15 21 * * *" # todos los días a las 00:15 Parámetros MSK: – nombre: fecha_inicio descripción: "Fecha de inicio (AAAA-MM-DD). Dejar vacío para auto" predeterminado: "" – nombre: fecha_final descripción: "Fecha de finalización (AAAA-MM-DD). Dejar vacío para auto" predeterminado: "" – nombre: meses_back descripción: "Meses para mirar hacia atrás (predeterminado: 5)" predeterminado: 5 alertas: habilitado: verdadero gravedad: advertencia
Luego, nuestro script Python analiza dag.yaml y dbt_project.yaml y utiliza la biblioteca Cosmos para generar un Airflow DAG completamente funcional. Esta es la única pieza de código Python en toda la configuración. Está escrito una vez y funciona para todos los proyectos dbt. Aquí está la parte clave:
def _build_dbt_project_dags(project_path: Path, environ: dict) -> list[DbtDag]: config_dict = yaml.safe_load(dag_config_path.read_text()) config = DagConfig.model_validate(config_dict) # Parámetros YAML → Parámetros de flujo de aire params = {} operator_vars = {} para parámetro en config.params: params[param.name] = Param( default=param.default if param.default no es Ninguno más "", descripción=param.description, ) operator_vars[param.name] = f"{{{{ params.{param.name} }}}}" # Cosmos crea el DAG a partir del proyecto dbt con DbtDag( dag_id=f"dbt_{project_path.name}", Schedule=config.schedule, params=params, project_config=ProjectConfig(dbt_project_path=project_path), perfil_config=ProfileConfig( perfil_name="default", target_name=project_name, perfil_mapping=TrinoLDAPProfileMapping( conn_id="trino_default", perfil_args={ "database": perfil_database, "schema": perfil_schema, }, ), ), operator_args={"vars": operator_vars}, ) as dag: # Crear esquema antes de ejecutar modelos create_schema = SQLExecuteQueryOperator( task_id="create_schema", conn_id="trino_default", sql=f"CREAR ESQUEMA SI NO EXISTE {profile_database}.{profile_schema} …", ) # Adjuntar a tareas raíz para Unique_id, _ en dag.dbt_graph.filtered_nodes.items(): task = dag.tasks_map[unique_id] si no es task.upstream_task_ids: create_schema >> tarea
Cosmos lee manifest.json del proyecto dbt, analiza el gráfico de dependencia del modelo y crea una tarea Airflow separada para cada modelo. Las dependencias de tareas se crean automáticamente en función de las llamadas ref() en SQL.
Cómo los analistas crean canales sin desarrolladores
Ahora, cuando un analista necesita un nuevo canal recurrente, puede armarlo en unos pocos pasos:
Paso 1. Cree una carpeta en el repositorio: dbt-projects/my_new_pipeline/.
Paso 2. Si se necesita la ingesta de datos externos, escriba una configuración YAML para dlt.
Paso 3. Escriba modelos SQL en la carpeta models/ y describa las fuentes en sources.yaml.
Paso 4. Cree dbt_project.yaml y dag.yaml.
Paso 5. Ingrese a Git, revise y fusione.
CI/CD crea el proyecto dbt y envía artefactos a S3. Airflow lee los archivos DAG desde allí, Cosmos analiza el proyecto dbt y genera el gráfico de tareas. Según lo previsto, dbt ejecuta los modelos en Trino en el orden correcto. El resultado final es un centro de datos actualizado en el almacén, al que se puede acceder a través de Superset.
Qué cambió después de la migración
Para que los analistas creen canalizaciones por su cuenta, deben comprender los conceptos de ref() y source(), la diferencia entre tabla y materialización incremental, y los conceptos básicos de Git. Realizamos algunos talleres internos y elaboramos guías paso a paso para cada tipo de tarea.
Por qué la nueva pila no reemplaza completamente a PySpark
Para aproximadamente el 10% de nuestras canalizaciones, PySpark sigue siendo la única opción, cuando una transformación simplemente no encaja en SQL. dbt admite macros Jinja, pero eso no sustituye al Python completo. Y sería deshonesto pasar por alto las limitaciones de las nuevas herramientas.
dlt + Delta: soporte de inserción experimental. Usamos el formato Delta en nuestra capa de almacenamiento. El conector Delta de dlt está marcado como experimental, por lo que la estrategia de fusión no funcionó de inmediato. Tuvimos que encontrar soluciones alternativas: en algunos casos usamos reemplazar en lugar de fusionar (sacrificando la incrementalidad) y en otros escribimos pasos de procesamiento personalizados.
La tolerancia limitada a fallas de Trino. Trino tiene un mecanismo de tolerancia a fallos, pero funciona escribiendo resultados intermedios en S3. En nuestros volúmenes de datos a escala de terabytes, esto no es práctico: la gran cantidad de operaciones S3 lo hace prohibitivamente costoso. Sin la tolerancia a fallos habilitada, si un trabajador de Trino deja de funcionar, toda la consulta falla. Spark, por el contrario, reinicia sólo la tarea fallida. Abordamos esto con reintentos a nivel de DAG y descomponiendo modelos pesados en cadenas de modelos intermedios.
UDF y lógica personalizada. En Spark, puedes escribir lógica personalizada en Python directamente dentro de la canalización, lo cual es muy conveniente. Con la nueva arquitectura, esto es mucho más difícil. dbt encima de Trino no ayuda: Jinja solo genera SQL y los modelos Python de dbt solo funcionan con Snowflake, Databricks y BigQuery. Puede escribir UDF en Trino, pero solo en Java, con toda la sobrecarga que eso implica: un repositorio separado, una canalización de compilación, implementación de JAR en todos los trabajadores. Entonces, cuando una transformación no encaja en SQL, terminas con un monstruo SQL que no se puede mantener o con un script independiente que rompe el linaje.
Qué sigue: pruebas, plantillas de modelos y capacitación
Mejores pruebas. Tuvimos pruebas sólidas en PySpark, pero la nueva arquitectura aún se está poniendo al día. Las versiones recientes de dbt introdujeron pruebas unitarias: ahora puede validar la lógica del modelo SQL con datos simulados sin acelerar el proceso completo. Queremos agregar pruebas dbt tanto a nivel de modelo como como una capa de monitoreo separada.
Plantillas reutilizables para patrones comunes. Muchos de nuestros modelos dbt se parecen. Una sola configuración podría describir una docena de modelos con el mismo patrón; solo difieren la tabla de origen y los filtros. Planeamos extraer la lógica compartida en macros dbt.
Ampliar la base de usuarios de la plataforma. Queremos que más ingenieros y analistas trabajen con datos de forma independiente. Estamos planificando sesiones periódicas de capacitación interna, documentación y guías de incorporación para que los nuevos usuarios puedan ponerse al día rápidamente y comenzar a construir sus propios modelos.
Si su equipo está atrapado en el mismo ciclo de "los analistas esperan a los desarrolladores", me encantaría saber cómo lo está resolviendo. Conéctate conmigo en LinkedIn y comparemos notas.
Todas las imágenes de este artículo son del autor a menos que se indique lo contrario.