Todo lo que hemos construido hasta ahora comparte una suposición cómoda: que los datos ya están en Google Cloud. El histórico venía de Cloud SQL, los eventos llegan por Pub/Sub, las exportaciones aterrizan en alpinashop-catalogo. Pero AlpinaShop, como cualquier empresa real, tiene datos que viven en otros sitios y llegan de otras maneras.

Son tres casos concretos, y ninguno es exótico:

  • El ERP del almacén es un MySQL 8 que corre en un servidor físico en las oficinas de Sabadell. Contiene el stock real, las entradas de mercancía, los proveedores y los costes de compra. Nadie va a migrarlo este año: funciona, está integrado con la báscula y con la impresora de etiquetas, y el proveedor cobra por cada cambio.
  • El transportista envía cada mes un CSV por correo electrónico con los envíos, las incidencias y los plazos de entrega reales. Separado por punto y coma, con fechas en formato dd/mm/aaaa, importes con coma decimal, nombres de ciudad en mayúsculas y sin acentos, y una columna DESTINATARIO que mezcla nombre y apellidos.
  • El proveedor de mochilas ofrece una API REST con su catálogo actualizado: referencias, precios de coste, disponibilidad y fichas técnicas.

Lucía necesita los tres. Necesita cruzar los costes de compra del ERP con las ventas de BigQuery para saber el margen real por producto. Necesita los plazos del transportista para explicar por qué caen las opiniones en ciertas zonas. Y necesita el catálogo del proveedor para detectar productos descatalogados que siguen publicados en la web.

Lucía sabe SQL. Sabe mucho SQL. Pero no ha escrito un pipeline de Apache Beam en su vida y no va a empezar ahora, y Dani tiene su propia lista de tareas. Si cada fichero nuevo requiere un desarrollo, esos datos no llegarán nunca.

Cloud Data Fusion existe para ese hueco: una herramienta de integración de datos visual, donde se construyen pipelines arrastrando bloques y se limpian datos viéndolos, sin escribir código. En esta lección la usarás para construir el pipeline del ERP, limpiarás el CSV del transportista con Wrangler, entenderás el linaje a nivel de campo que es su mayor virtud, montarás la replicación con captura de cambios desde MySQL, y —esto es igual de importante— aprenderás cuándo no usarla, porque tiene un coste por hora que condiciona toda la decisión.

Contenido

  1. El problema de los datos que no nacen en Google Cloud
  2. Qué es un ETL/ELT visual y a quién sirve
  3. Qué es Data Fusion: CDAP gestionado
  4. Ediciones, coste por hora y sus consecuencias
  5. Crear la instancia y entender qué se ha creado
  6. El Studio: orígenes, transformaciones y destinos
  7. Wrangler: limpiar el CSV del transportista viendo los datos
  8. El pipeline real: MySQL del ERP a alpinashop_analitica
  9. Conectores y plugins del Hub
  10. Desplegar, ejecutar y programar
  11. Qué ocurre por debajo: Dataproc efímero
  12. Linaje de datos a nivel de campo
  13. Replicación con CDC desde el MySQL del ERP
  14. Cuándo Data Fusion, cuándo Dataflow, cuándo Datastream, cuándo bq load
  15. La decisión razonada de AlpinaShop

  1. El problema de los datos que no nacen en Google Cloud

Se le llama integración de datos, y es el trabajo menos glamuroso y más consumidor de tiempo de cualquier proyecto analítico. Las encuestas del sector lo repiten sin variación: entre el 60 % y el 80 % del esfuerzo de un proyecto de datos se va en conseguir que los datos lleguen, no en analizarlos.

Los problemas concretos, con los ejemplos de AlpinaShop:

Problema Caso en AlpinaShop
Conectividad El MySQL está en Sabadell, detrás de un firewall, sin IP pública
Formatos CSV con ;, fechas dd/mm/aaaa, decimales con coma, sin cabecera fiable
Calidad Ciudades en mayúsculas sin acentos, nulos escritos como -, espacios de sobra
Esquemas cambiantes El transportista añadió una columna en enero sin avisar
Frecuencia dispar El ERP cambia continuamente; el CSV llega mensual; la API se consulta a demanda
Volumen El histórico del ERP son 12 millones de movimientos de stock
Trazabilidad Nadie sabe de dónde salió el campo coste_unitario de un informe de 2025

Cada uno de estos problemas es resoluble programando. El asunto es que resolverlos programando para cada origen, y mantener ese código cuando el transportista cambia el formato, es un trabajo continuo que una pyme de 40 personas no puede permitirse dedicar a Dani.

  1. Qué es un ETL/ELT visual y a quién sirve

ETL significa extraer, transformar, cargar: se saca el dato del origen, se transforma fuera, y se carga ya limpio en el destino. ELT invierte los dos últimos: se carga el dato en crudo y se transforma dentro del destino, aprovechando su potencia.

Enfoque Dónde se transforma Ventaja Inconveniente
ETL En un motor intermedio El destino solo recibe datos limpios Ese motor hay que dimensionarlo y pagarlo
ELT En el destino (BigQuery) Aprovecha su motor; conservas el crudo El destino guarda datos sucios; coste de consulta

Con almacenes modernos como BigQuery, la tendencia clara es ELT: cargar crudo y transformar con SQL. Pero hay una parte que sigue siendo ETL inevitablemente: sacar el dato del origen y llevarlo hasta la puerta. Eso es la extracción, y es exactamente donde Data Fusion aporta.

Un ETL visual es una herramienta donde el pipeline se construye con un lienzo y bloques en lugar de con código. ¿A quién sirve?

Sirve muy bien a:

  • Analistas como Lucía, que conocen el negocio y los datos pero no programan pipelines distribuidos.
  • Equipos pequeños sin ingenieros de datos dedicados.
  • Integraciones estándar: leer una tabla, limpiar unos campos, escribir en otra tabla.
  • Organizaciones que necesitan trazabilidad documentada de dónde sale cada dato (auditorías, cumplimiento).

Sirve mal a:

  • Lógica de negocio compleja con condiciones anidadas y estado.
  • Equipos que ya tienen ingenieros y control de versiones maduro: un pipeline visual se versiona peor que un fichero .py.
  • Streaming con ventanas y tiempo del evento: para eso está Beam.
  • Presupuestos ajustados con uso esporádico, por el motivo del apartado 4.

Y una advertencia honesta que conviene hacer desde el principio: sin código no significa sin conocimiento. Para construir un pipeline en Data Fusion hay que entender esquemas, tipos, uniones, claves y particiones exactamente igual que programándolo. Lo que se ahorra es la sintaxis y la infraestructura, no el pensamiento.

  1. Qué es Data Fusion: CDAP gestionado

Cloud Data Fusion es la versión gestionada de CDAP (Cask Data Application Platform), una plataforma open source de integración de datos que Google adquirió y ofrece como servicio.

Sus componentes:

Componente Qué hace
Studio El lienzo visual donde se dibuja el pipeline
Wrangler Explorador y limpiador interactivo de datos
Hub Catálogo de plugins, conectores y pipelines de ejemplo
Metadatos y linaje Registro automático de qué campo viene de dónde
Replicación Módulo de CDC para copiar bases de datos en continuo
Motor de ejecución Genera y ejecuta el trabajo en Dataproc

La última fila es la más importante para entender el producto: Data Fusion no ejecuta nada por sí mismo. Traduce el pipeline visual a un trabajo de Spark o MapReduce y lo ejecuta en un clúster de Dataproc que crea al vuelo. Eso explica su comportamiento, su tiempo de arranque y buena parte de su coste, y lo desarrollaremos en el apartado 11.

Ser CDAP gestionado tiene otra consecuencia relevante: los pipelines son portables. Un pipeline exportado de Data Fusion se puede importar en un CDAP autogestionado en cualquier sitio. No es un formato propietario cerrado.

  1. Ediciones, coste por hora y sus consecuencias

Aquí está la característica que condiciona todas las decisiones sobre este producto, y hay que ponerla por delante en lugar de esconderla al final.

Edición Para qué Coste aproximado por hora de instancia Notas
Developer Pruebas y desarrollo ~0,35 $ Sin alta disponibilidad, capacidad limitada
Basic Producción sencilla ~1,80 $ Incluye 120 horas gratuitas al mes por cuenta
Enterprise Producción exigente ~4,20 $ Alta disponibilidad, más concurrencia, linaje completo, CDC

Verifica los precios vigentes en la documentación oficial; lo que importa es el modelo, no la cifra exacta.

Y el modelo es este: se paga por hora de instancia existente, se use o no. No por pipeline ejecutado, ni por dato procesado. La instancia es un entorno que está encendido.

Hagamos el cálculo para AlpinaShop, porque es la conversación que hay que tener con dirección:

Escenario Horas/mes Coste aproximado
Instancia Basic permanente 730 ~1.310 $ (menos 120 h gratis: ~1.100 $)
Instancia Enterprise permanente 730 ~3.070 $
Instancia Developer permanente 730 ~255 $
Basic encendida 4 h al día 120 0 $ (dentro de las 120 h gratuitas)

Y a eso hay que sumarle el coste del clúster de Dataproc que se levanta para ejecutar cada pipeline, que en la práctica suele ser menor pero no es cero.

La consecuencia es directa y hay que decirla sin adornos: una instancia de Data Fusion permanente cuesta más al mes que todo el resto de la plataforma de datos de AlpinaShop junta. BigQuery cuesta unos pocos euros, Pub/Sub céntimos, Dataflow por lotes céntimos. Data Fusion Basic permanente costaría más de mil euros.

Eso no invalida el producto: para una empresa con veinte orígenes de datos y un equipo de integración, mil euros al mes es una ganga frente a los salarios que ahorra. Pero para una pyme con tres orígenes, la ecuación no sale, y hay que decirlo.

La estrategia que sí funciona en una pyme es tratar la instancia como los clústeres de Dataproc de 04-03: efímera. Se enciende para desarrollar, se apaga al terminar; y los pipelines ya desplegados se ejecutan igual, porque el trabajo lo hace Dataproc. Volveremos a ello en el apartado 15.

  1. Crear la instancia y entender qué se ha creado

gcloud config set project alpinashop-datos
gcloud services enable datafusion.googleapis.com

# Instancia de desarrollo: la mas barata para aprender y construir
gcloud data-fusion instances create alpinashop-fusion \
  --location=europe-west1 \
  --type=DEVELOPER \
  --enable-stackdriver-logging \
  --enable-stackdriver-monitoring \
  --labels=entorno=desarrollo,equipo=datos,centro-coste=analitica

La creación tarda entre 15 y 25 minutos. No es un error: se está aprovisionando un entorno completo de CDAP en un proyecto gestionado por Google.

Ese detalle del "proyecto gestionado" importa para los permisos. Data Fusion crea la instancia en un proyecto propio de Google (el tenant project) y actúa sobre el tuyo mediante un agente de servicio, al que hay que conceder permisos explícitamente:

PROY_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_FUSION="service-${PROY_NUM}@gcp-sa-datafusion.iam.gserviceaccount.com"

# El agente necesita poder actuar sobre el proyecto
gcloud projects add-iam-policy-binding alpinashop-datos \
  --member="serviceAccount:${SA_FUSION}" \
  --role="roles/datafusion.serviceAgent"

# Cuenta de servicio que usaran los clusteres de Dataproc que ejecutan los pipelines
gcloud iam service-accounts create sa-fusion-pipelines \
  --display-name="Ejecucion de pipelines de Data Fusion"

SA_PIPE="[email protected]"

for ROL in roles/dataproc.worker roles/bigquery.dataEditor \
           roles/bigquery.jobUser roles/storage.objectAdmin; do
  gcloud projects add-iam-policy-binding alpinashop-datos \
    --member="serviceAccount:${SA_PIPE}" --role="$ROL"
done

# El agente de Data Fusion debe poder usar esa cuenta
gcloud iam service-accounts add-iam-policy-binding "$SA_PIPE" \
  --member="serviceAccount:${SA_FUSION}" \
  --role="roles/iam.serviceAccountUser"

Acceso a la interfaz:

gcloud data-fusion instances describe alpinashop-fusion \
  --location=europe-west1 --format="value(apiEndpoint, serviceEndpoint, state)"

Y la conectividad con Sabadell, que es el requisito real para leer del ERP. Hay tres opciones, en orden de preferencia:

Opción Cómo funciona Valoración
Cloud VPN o Interconnect Túnel entre alpinashop-vpc y la red de las oficinas La correcta. El MySQL nunca se expone a internet
IP pública con lista blanca y TLS Se abre el puerto solo a rangos concretos Aceptable a regañadientes; superficie de ataque innecesaria
Exportar a fichero y subirlo Un script del ERP vuelca CSV a Cloud Storage Sencillo, pero no permite CDC ni datos frescos

Para AlpinaShop, Marta monta un túnel de Cloud VPN con enrutamiento hacia la subred del ERP, coherente con todo lo visto en 03-01: el MySQL sigue sin IP pública y el tráfico va cifrado. La configuración detallada de la conectividad híbrida es territorio de 07-03.

  1. El Studio: orígenes, transformaciones y destinos

El Studio es un lienzo. A la izquierda hay una paleta de bloques agrupados por categoría; se arrastran al lienzo y se conectan con flechas que representan el flujo de los datos.

Los bloques se clasifican así:

Categoría Qué hace Ejemplos
Source Lee datos Database (JDBC), BigQuery, GCS, Salesforce, HTTP, Kafka
Transform Modifica registros Wrangler, JavaScript, Python, Projection, Encoder
Analytics Agrega y une Group By, Joiner, Deduplicate, Distinct, Row Denormalizer
Conditions and Actions Control de flujo y acciones Condition, Email, BigQuery Execute, Database Execute
Sink Escribe datos BigQuery, GCS, Database, Spanner, Pub/Sub
Error Handlers Recogen registros rechazados Error Collector

Cada bloque tiene un panel de configuración: cadena de conexión, tabla, esquema de salida, opciones específicas. Y cada bloque declara un esquema de salida —la lista de campos con sus tipos— que se propaga al siguiente. Esa propagación es lo que permite el linaje del apartado 12.

El pipeline que vamos a construir, dibujado:

flowchart LR
    S1["Source: Database<br/>MySQL ERP Sabadell<br/>tabla movimientos_stock"]
    S2["Source: GCS<br/>CSV del transportista"]
    W1["Transform: Wrangler<br/>limpieza y tipos"]
    W2["Transform: Wrangler<br/>fechas, decimales, nombres"]
    J["Analytics: Joiner<br/>por sku"]
    G["Analytics: Group By<br/>coste medio por sku"]
    K1["Sink: BigQuery<br/>alpinashop_analitica.stock_erp"]
    K2["Sink: BigQuery<br/>alpinashop_analitica.envios"]
    E["Error Collector<br/>-> GCS cuarentena"]

    S1 --> W1 --> G --> K1
    S2 --> W2 --> J
    W1 --> J
    J --> K2
    W2 -.rechazos.-> E

Fíjate en dos cosas del diagrama, porque reproducen las buenas prácticas de las lecciones anteriores:

  • La rama de rechazos existe también aquí. El bloque Error Collector recoge los registros que un Wrangler no pudo procesar y los envía a un destino de cuarentena, exactamente igual que las salidas etiquetadas de Beam en 04-02.
  • Un mismo origen alimenta dos ramas. El lienzo es un grafo dirigido, no una línea, igual que el grafo de Beam.

  1. Wrangler: limpiar el CSV del transportista viendo los datos

Wrangler es la pieza que justifica por sí sola aprender Data Fusion. Es un entorno interactivo donde cargas una muestra de los datos, la ves en forma de tabla, y aplicas transformaciones —llamadas directivas— viendo el efecto inmediatamente.

Partimos del CSV real del transportista, tal como llega:

NUM_ENVIO;FECHA_ENTREGA;DESTINATARIO;CIUDAD;CP;PESO;IMPORTE;INCIDENCIA;PEDIDO
ENV0098211;14/03/2026;GARCIA LOPEZ, MARIA;BARCELONA;08013;2,450;4,90;-;PED-2026-0042
ENV0098212;15/03/2026;  martinez ruiz, juan  ;VALENCIA;46001;1,200;3,50;RETRASO 24H;PED-2026-0043
ENV0098213;-;FERNANDEZ SANZ, ANA;MADRID;28004;5,000;12,00;DIRECCION INCORRECTA;PED-2026-0044

Los problemas saltan a la vista: separador ;, fechas europeas, decimales con coma, nulos como -, espacios sobrantes, mayúsculas inconsistentes, nombre y apellidos juntos, y una entrega sin fecha.

Las directivas de Wrangler se escriben una por línea y se aplican en orden. Se pueden generar desde el menú contextual de cada columna, pero es mucho más rápido escribirlas:

-- 1) Partir la linea por el separador y nombrar las columnas
parse-as-csv :body ';' true
drop :body

-- 2) Limpiar espacios sobrantes en todas las columnas de texto
trim :DESTINATARIO
trim :CIUDAD
trim :INCIDENCIA

-- 3) Nulos: el transportista escribe '-' donde no hay dato
find-and-replace :FECHA_ENTREGA s/^-$//g
find-and-replace :INCIDENCIA s/^-$//g
set-column :INCIDENCIA (INCIDENCIA == null || INCIDENCIA.isEmpty()) ? null : INCIDENCIA

-- 4) Fechas europeas -> tipo fecha real
parse-as-simple-date :FECHA_ENTREGA dd/MM/yyyy
format-date :FECHA_ENTREGA yyyy-MM-dd

-- 5) Decimales con coma -> punto, y luego a numero
find-and-replace :PESO s/,/./g
find-and-replace :IMPORTE s/,/./g
set-type :PESO double
set-type :IMPORTE double

-- 6) Partir 'APELLIDOS, NOMBRE' en dos columnas
split-to-columns :DESTINATARIO ,
rename :DESTINATARIO_1 apellidos
rename :DESTINATARIO_2 nombre
trim :apellidos
trim :nombre

-- 7) Normalizar el formato de los nombres propios
titlecase :apellidos
titlecase :nombre
uppercase :CIUDAD

-- 8) Columna derivada: hubo incidencia si el campo no esta vacio
set-column :hubo_incidencia (INCIDENCIA != null)

-- 9) Nombres de columna en minusculas y coherentes con el resto del almacen
rename :NUM_ENVIO envio_id
rename :FECHA_ENTREGA fecha_entrega
rename :CIUDAD ciudad
rename :CP codigo_postal
rename :PESO peso_kg
rename :IMPORTE importe_envio_eur
rename :INCIDENCIA incidencia
rename :PEDIDO pedido_id

-- 10) Descartar filas sin identificador de pedido: no sirven para nada
filter-rows-on condition-false pedido_id != null && !pedido_id.isEmpty()

-- 11) Quitar la columna de nombre y apellidos por MINIMIZACION DE DATOS
drop :apellidos
drop :nombre

Las directivas más útiles, agrupadas:

Familia Directivas Para qué
Parseo parse-as-csv, parse-as-json, parse-as-xml, parse-as-fixed-length Convertir texto crudo en columnas
Fechas parse-as-simple-date, format-date, parse-as-datetime Formatos regionales
Texto trim, uppercase, lowercase, titlecase, cleanse-column-names Normalización
Búsqueda find-and-replace, extract-regex-groups, split-to-columns Expresiones regulares
Tipos set-type, fill-null-or-empty Conversión y nulos
Filtros filter-rows-on, filter-row-if-matched Descartar filas
Columnas rename, drop, keep, set-column, merge Estructura
Calidad send-to-error Desviar registros inválidos

Los pasos 10 y 11 del ejemplo merecen comentario.

El paso 10 usa filter-rows-on, que descarta silenciosamente. Si prefieres conservar lo descartado para revisarlo —y normalmente deberías—, se usa send-to-error, que envía el registro al Error Collector del pipeline en lugar de tirarlo:

send-to-error pedido_id == null || pedido_id.isEmpty()

El paso 11 es una decisión de cumplimiento, no de limpieza:

Aviso de RGPD. El CSV del transportista contiene el nombre y apellidos del destinatario, que es un dato personal identificativo. Para el análisis de plazos de entrega e incidencias, ese dato no aporta absolutamente nada: basta con el código postal y el identificador de pedido. Eliminarlo en la fase de limpieza, antes de que llegue al almacén analítico, es minimización de datos por diseño y es la práctica correcta. Cualquier tratamiento de datos personales reales debe ser revisado por un profesional de compliance o el DPO antes de pasar a producción. Todos los datos de este curso son ficticios.

La ventaja de Wrangler frente a escribir esto en Python o SQL no es la potencia —Python puede hacer todo esto y más—, sino el ciclo de retroalimentación: aplicas una directiva y ves inmediatamente el efecto sobre 100 filas reales. Cuando parse-as-simple-date falla porque hay una fecha con formato distinto en la fila 47, lo ves en el momento, no cuando el pipeline lleve veinte minutos en producción.

  1. El pipeline real: MySQL del ERP a alpinashop_analitica

Vamos con la integración que de verdad quiere Lucía: traer los costes de compra y el stock del ERP para poder calcular el margen real por producto.

Bloque 1 — Source: Database. Se configura con:

  • Plugin type: Database (genérico JDBC).
  • JDBC driver: el conector de MySQL, que hay que subir previamente desde el Hub o como plugin propio.
  • Connection string: jdbc:mysql://10.20.0.15:3306/erp_almacen (IP privada alcanzable por la VPN).
  • Import Query: la consulta que extrae los datos.
-- Consulta de extraccion del ERP.
-- $CONDITIONS es OBLIGATORIO si se usa lectura particionada: Data Fusion
-- lo sustituye por un rango distinto en cada tarea paralela.
SELECT
  m.sku,
  m.fecha_movimiento,
  m.tipo_movimiento,
  m.cantidad,
  m.coste_unitario,
  m.proveedor_id,
  p.nombre AS proveedor_nombre,
  m.almacen
FROM movimientos_stock m
LEFT JOIN proveedores p ON p.id = m.proveedor_id
WHERE m.fecha_movimiento >= '${fecha_desde}'
  AND m.fecha_movimiento <  '${fecha_hasta}'
  AND $CONDITIONS

Dos detalles fundamentales de esta configuración:

  • ${fecha_desde} y ${fecha_hasta} son macros de Data Fusion: argumentos en tiempo de ejecución. Permiten que el mismo pipeline sirva para la carga inicial completa y para la incremental diaria, sin duplicarlo. El orquestador de 04-06 les pasará los valores.
  • $CONDITIONS con numSplits: si configuras Split-By Field Name = sku y Number of Splits = 4, Data Fusion lanza cuatro consultas en paralelo sobre rangos distintos. Esto acelera muchísimo la carga inicial, pero pone cuatro veces más carga sobre el MySQL del ERP. Con una base de datos de producción que además atiende a la báscula del almacén, hay que ser prudente: uno o dos splits, y ejecutar de madrugada.

Bloque 2 — Transform: Wrangler. Limpieza específica del ERP:

-- El ERP mezcla mayusculas y minusculas en los SKU
uppercase :sku
trim :sku

-- Codigos de movimiento internos -> etiquetas legibles
set-column :tipo_movimiento (tipo_movimiento == 'E' ? 'entrada' : (tipo_movimiento == 'S' ? 'salida' : 'ajuste'))

-- El ERP guarda los costes en centimos como entero
set-column :coste_unitario_eur coste_unitario / 100.0
drop :coste_unitario

-- Descartar movimientos sin SKU: son ajustes contables, no de stock
send-to-error sku == null || sku.isEmpty()

-- Marca de cuando se extrajo, para poder auditar cargas
set-column :cargado_en datetime:CurrentDateTime()

Bloque 3 — Analytics: Group By. Coste medio ponderado por SKU:

  • Group by fields: sku
  • Aggregates:
    • Sum(cantidad) → unidades_compradas
    • Avg(coste_unitario_eur) → coste_medio_eur
    • Min(fecha_movimiento) → primera_compra
    • Max(fecha_movimiento) → ultima_compra
    • Count(*) → num_movimientos

Bloque 4 — Sink: BigQuery.

  • Dataset: alpinashop_analitica
  • Table: costes_producto_erp
  • Operation: Upsert con clave sku (actualiza si existe, inserta si no)
  • Truncate Table: desactivado
  • Service Account: [email protected]
  • Location: europe-west1

La opción Upsert es la que hace el pipeline reejecutable sin duplicar. Es el equivalente al MERGE de SQL, y es lo que permite lanzar el pipeline dos veces sin estropear nada, la misma propiedad de idempotencia que perseguíamos en 04-04.

Bloque 5 — Error Collector → Sink GCS. Los registros desviados con send-to-error van a gs://alpinashop-datalake/cuarentena/erp/, en formato JSON, con el motivo del rechazo.

Y con esos datos ya en BigQuery, Lucía puede por fin responder a su pregunta con SQL:

-- Margen real por producto: precio de venta medio contra coste del ERP
SELECT
  pr.categoria,
  pr.sku,
  pr.nombre,
  ROUND(AVG(l.precio_unitario), 2)                      AS precio_medio_venta,
  ROUND(c.coste_medio_eur, 2)                           AS coste_medio,
  ROUND(AVG(l.precio_unitario) - c.coste_medio_eur, 2)  AS margen_eur,
  ROUND(100 * (AVG(l.precio_unitario) - c.coste_medio_eur)
        / NULLIF(AVG(l.precio_unitario), 0), 1)         AS margen_pct,
  SUM(l.cantidad)                                       AS unidades_vendidas
FROM `alpinashop-datos.alpinashop_analitica.lineas_pedido`      AS l
JOIN `alpinashop-datos.alpinashop_analitica.productos`          AS pr USING (sku)
JOIN `alpinashop-datos.alpinashop_analitica.costes_producto_erp` AS c  USING (sku)
WHERE l.fecha_pedido >= DATE '2026-01-01'
GROUP BY pr.categoria, pr.sku, pr.nombre, c.coste_medio_eur
HAVING unidades_vendidas > 10
ORDER BY margen_pct ASC;      -- los peores margenes primero: eso es lo accionable

Esa consulta, que ordena por el margen más bajo, es exactamente el tipo de resultado que cambia decisiones: productos que se venden mucho y dejan poco. No era posible antes de esta lección porque el coste vivía en un servidor de Sabadell.

  1. Conectores y plugins del Hub

El Hub es el catálogo desde el que se instalan plugins en la instancia. Categorías principales:

Tipo Ejemplos disponibles
Bases de datos MySQL, PostgreSQL, SQL Server, Oracle, DB2, Teradata, MongoDB
Google Cloud BigQuery, GCS, Spanner, Bigtable, Pub/Sub, Datastore
SaaS Salesforce, SAP, ServiceNow, Marketo, Zendesk, Google Analytics
Ficheros y protocolos HTTP, FTP/SFTP, Amazon S3, Azure Blob, Excel, XML
Transformaciones Wrangler, JavaScript, Python Evaluator, XML Parser, Validator
Analytics Joiner, Group By, Deduplicate, Pivot, Window Aggregation

Para el tercer origen de AlpinaShop, la API del proveedor de mochilas, se usa el plugin HTTP:

  • URL: https://api.proveedor-montana.example/v2/catalogo
  • HTTP Method: GET
  • Headers: Authorization: Bearer ${api_token} — donde ${api_token} es una macro que se resuelve desde Secret Manager (03-06), nunca escrita en el pipeline.
  • Format: json
  • JSON/XML Result Path: $.productos — la ruta dentro de la respuesta donde está el array.
  • Pagination Type: Link Header o Increment an Index, según lo que soporte la API.

La paginación es lo que más se olvida: sin configurarla, se traen solo los primeros N resultados y nadie se da cuenta hasta que faltan productos.

Si un origen no tiene plugin, quedan tres salidas: escribirlo en Java (CDAP es extensible), usar el Python Evaluator para transformaciones puntuales, o —lo más razonable— exportar el dato a Cloud Storage con un script y leerlo desde ahí. La última suele ser la respuesta correcta para un origen exótico y de bajo volumen.

  1. Desplegar, ejecutar y programar

Un pipeline en Data Fusion tiene dos estados: borrador (se edita, se previsualiza) y desplegado (inmutable, ejecutable, versionado).

El flujo de trabajo:

  1. Preview. Ejecuta el pipeline con una muestra pequeña sin escribir en el destino. Muestra los datos que salen de cada bloque. Es el equivalente al DirectRunner de Beam y hay que usarlo siempre antes de desplegar.
  2. Deploy. Congela el pipeline con un número de versión. Para cambiarlo, se crea una versión nueva.
  3. Run. Ejecuta. Se pueden pasar argumentos en tiempo de ejecución (las macros).
  4. Schedule. Programa ejecuciones periódicas con expresión cron.

Todo esto también se hace por API REST, que es como lo invocará el orquestador de 04-06:

INSTANCIA=$(gcloud data-fusion instances describe alpinashop-fusion \
  --location=europe-west1 --format="value(apiEndpoint)")
TOKEN=$(gcloud auth print-access-token)

# Ejecutar el pipeline con macros de fecha
curl -X POST \
  -H "Authorization: Bearer ${TOKEN}" \
  -H "Content-Type: application/json" \
  "${INSTANCIA}/v3/namespaces/default/apps/erp-costes-a-bigquery/workflows/DataPipelineWorkflow/start" \
  -d '{
    "fecha_desde": "2026-03-01",
    "fecha_hasta": "2026-04-01",
    "system.profile.name": "alpinashop-perfil-computo"
  }'

# Consultar el estado de las ejecuciones
curl -H "Authorization: Bearer ${TOKEN}" \
  "${INSTANCIA}/v3/namespaces/default/apps/erp-costes-a-bigquery/workflows/DataPipelineWorkflow/runs"

Los perfiles de cómputo (compute profiles) definen cómo será el clúster de Dataproc que ejecute el pipeline: número de workers, tipo de máquina, red, cuenta de servicio. Es donde se controla el coste de ejecución:

Perfil "alpinashop-perfil-computo"
  Provisioner       : Dataproc
  Region            : europe-west1
  Master            : n2-standard-2, 1 nodo, disco 100 GB
  Workers           : n2-standard-2, 2 nodos, disco 100 GB
  Network           : alpinashop-vpc / sn-datos-euw1
  Internal IP only  : true
  Service Account   : [email protected]
  Image version     : 2.2-debian12
  Idle TTL          : 10 minutos

Ese perfil aplica exactamente lo aprendido en 04-03: red privada sin IP pública, cuenta de servicio propia, versión de imagen fijada y autodestrucción por inactividad.

  1. Qué ocurre por debajo: Dataproc efímero

Cuando pulsas Run, esto es lo que pasa realmente:

sequenceDiagram
    participant U as Lucia
    participant DF as Data Fusion
    participant DP as Dataproc
    participant BQ as BigQuery

    U->>DF: Run del pipeline
    DF->>DF: traduce el grafo visual a un job de Spark
    DF->>DP: crea un cluster efimero (2-5 min)
    DP->>DP: ejecuta el job de Spark
    DP->>BQ: escribe el resultado
    DP-->>DF: fin del job
    DF->>DP: destruye el cluster (Idle TTL)
    DF-->>U: estado SUCCEEDED + linaje registrado

Consecuencias muy prácticas de saber esto:

El arranque tarda de 2 a 5 minutos. No es lentitud de Data Fusion: es el tiempo de crear el clúster. Por eso Data Fusion no sirve para nada que necesite latencia baja. Un pipeline que mueve 200 filas tarda casi lo mismo que uno que mueve 20 millones, porque el coste dominante es el arranque.

El coste real es la suma de dos cosas: las horas de instancia (siempre) más las horas de Dataproc (por ejecución). Un pipeline diario de 8 minutos con 3 nodos son unos céntimos de Dataproc, pero la instancia sigue facturando 24 horas al día.

Puedes reutilizar un clúster. Configurando un perfil que apunte a un clúster existente en lugar de crear uno, se elimina el tiempo de arranque. Tiene sentido cuando se ejecutan muchos pipelines seguidos, por ejemplo en una ventana nocturna: se crea el clúster, se lanzan diez pipelines encadenados, se destruye.

Los errores de Spark aparecen en los logs de Dataproc. Cuando un pipeline falla por memoria (OutOfMemoryError en un executor) o por sesgo de datos, el diagnóstico es exactamente el de 04-03: mirar la interfaz de Spark y las etapas. Data Fusion no te oculta esa realidad, solo la envuelve.

  1. Linaje de datos a nivel de campo

Esta es, para muchas organizaciones, la razón principal para pagar Data Fusion.

El linaje responde a dos preguntas que suenan triviales y que casi ninguna empresa sabe contestar:

  • "¿De dónde sale exactamente este campo del informe?"
  • "Si cambio esta columna del ERP, ¿qué se rompe?"

Data Fusion registra el linaje automáticamente al ejecutar cada pipeline, a dos niveles:

Linaje de conjunto de datos: qué orígenes alimentan qué destinos, con qué pipeline y cuándo.

Linaje de campo: qué columna concreta del origen produce qué columna del destino, y qué operaciones ha sufrido por el camino.

flowchart LR
    A["erp_almacen.movimientos_stock<br/>coste_unitario (INT, centimos)"]
    B["Wrangler<br/>coste_unitario / 100.0"]
    C["Group By<br/>Avg()"]
    D["alpinashop_analitica.costes_producto_erp<br/>coste_medio_eur (DOUBLE)"]
    A --> B --> C --> D

Ese diagrama no lo dibuja nadie: lo genera Data Fusion a partir de la ejecución. Y responde a la pregunta que en muchas empresas cuesta días de arqueología: si alguien pregunta por qué el margen del informe de dirección sale raro, el linaje muestra que coste_medio_eur viene de una división entre 100 y un promedio no ponderado —que, dicho sea de paso, es una decisión discutible que el linaje deja a la vista.

Por qué importa tanto:

Escenario Sin linaje Con linaje
Auditoría "Creemos que viene del ERP" Trazabilidad documentada y fechada
Cambio en el origen Se despliega y se ve qué se rompe Se sabe de antemano qué pipelines dependen
Dato sospechoso Días de investigación Un clic hasta el origen
RGPD: dónde está un dato personal Búsqueda manual por todas partes Se rastrea el campo por todos los destinos
Baja de un sistema Miedo a apagarlo Se ve exactamente qué consume de él

La última fila es especialmente valiosa en una migración: saber qué depende del ERP antes de tocarlo. Y la penúltima conecta directamente con lo que veremos en 04-07, donde Dataplex consolida el linaje de toda la plataforma, no solo el de Data Fusion.

  1. Replicación con CDC desde el MySQL del ERP

Un pipeline por lotes que se ejecuta cada noche deja los datos con hasta 24 horas de retraso, y además carga el ERP con una consulta pesada. La alternativa es la captura de datos de cambio (Change Data Capture, CDC): en lugar de consultar la tabla, se lee el registro de transacciones de la base de datos y se replican los cambios según ocurren.

flowchart LR
    M["MySQL ERP Sabadell<br/>binlog"]
    R["Data Fusion Replication<br/>lee el binlog"]
    S["Tabla de staging<br/>en BigQuery"]
    B["Tabla destino<br/>alpinashop_analitica.stock_erp"]

    M -->|INSERT/UPDATE/DELETE| R
    R -->|eventos de cambio| S
    S -->|MERGE periodico| B

Preparación en el MySQL del ERP:

-- Requisitos en el servidor de origen (los aplica el administrador del ERP)
-- En my.cnf:
--   server-id        = 1
--   log_bin          = mysql-bin
--   binlog_format    = ROW          <-- IMPRESCINDIBLE: ROW, no STATEMENT
--   binlog_row_image = FULL
--   expire_logs_days = 7

CREATE USER 'cdc_datafusion'@'%' IDENTIFIED BY 'contrasena-desde-secret-manager';
GRANT SELECT, RELOAD, SHOW DATABASES,
      REPLICATION SLAVE, REPLICATION CLIENT
  ON *.* TO 'cdc_datafusion'@'%';
FLUSH PRIVILEGES;

binlog_format = ROW es innegociable: con STATEMENT, el binlog guarda las sentencias SQL, no las filas resultantes, y la replicación no puede reconstruir el estado. Es el primer requisito que hay que verificar con el proveedor del ERP.

En Data Fusion se usa el módulo Replication (edición Enterprise):

  1. Origen: MySQL, con la conexión y el usuario cdc_datafusion.
  2. Selección de tablas: movimientos_stock, proveedores, articulos.
  3. Destino: BigQuery, dataset alpinashop_analitica, con un prefijo de staging.
  4. Evaluación de la fuente: comprueba permisos y configuración antes de empezar.
  5. Instantánea inicial + streaming continuo.

Data Fusion carga primero una foto completa de las tablas y después aplica los cambios en continuo, escribiendo en tablas de staging y ejecutando un MERGE periódico contra las tablas finales.

Consideraciones importantes sobre CDC, que hay que valorar antes de decidirse:

Aspecto Realidad
Latencia Segundos o pocos minutos, frente a 24 horas del lote
Carga en el origen Mucho menor: leer el binlog no ejecuta consultas
Borrados Se capturan, cosa que una carga incremental por fecha no hace
Coste Requiere edición Enterprise y ejecución continua: es caro
Requisitos Configuración del servidor de origen; a veces el proveedor no la permite
Cambios de esquema Se manejan, pero requieren atención

Ese punto de los borrados es el argumento técnico más fuerte a favor de CDC. Un pipeline incremental que lee WHERE fecha_modificacion > X nunca se entera de que una fila se ha eliminado, y el almacén analítico acumula registros fantasma indefinidamente. CDC sí ve el DELETE.

La alternativa a considerar es Datastream, un servicio de Google dedicado exclusivamente a CDC, más simple y más barato que activar Enterprise en Data Fusion solo para esto. Está en la tabla del apartado siguiente.

  1. Cuándo Data Fusion, cuándo Dataflow, cuándo Datastream, cuándo bq load

La tabla honesta, que es lo que hay que llevarse de esta lección:

Criterio bq load Datastream Data Fusion Dataflow
Caso típico Fichero listo → tabla Réplica de BD con CDC Integrar orígenes variados con limpieza Transformación a medida, streaming
Código Un comando Ninguno Ninguno (visual) Python o Java
Perfil necesario Cualquiera Cualquiera + DBA del origen Analista Ingeniero de datos
Coste fijo Cero Por GB procesado Por hora de instancia Cero en lote
Coste variable Gratis Bajo Dataproc por ejecución Trabajadores por hora
Latencia mínima Minutos Segundos 2-5 min de arranque Segundos (streaming)
Transformaciones Ninguna Ninguna (replica tal cual) Muchas, visuales Ilimitadas
Streaming real No Sí (replicación) Limitado Sí, su punto fuerte
Conectores externos No Bases de datos Muchísimos Los que programes
Linaje No No Sí, por campo Vía Dataplex
Versionado en Git Trivial N/A Incómodo (JSON exportado) Trivial

Y traducido a reglas de decisión:

Usa bq load cuando el dato ya está en Cloud Storage con la forma correcta. Es gratis y no hay nada que mantener. Empieza siempre preguntándote si esto basta. La mitad de los pipelines que existen en el mundo sobran.

Usa Datastream cuando el objetivo sea replicar una base de datos completa en BigQuery con baja latencia y captura de borrados, sin transformar nada. Es más simple y barato que Data Fusion Enterprise para esa tarea concreta, y admite MySQL, PostgreSQL, Oracle y SQL Server.

Usa Data Fusion cuando haya varios orígenes heterogéneos, hagan falta transformaciones y limpieza no triviales, el equipo no tenga perfil de programación, y el linaje sea un requisito real (auditoría, cumplimiento). Y cuando el volumen de trabajo justifique el coste de la instancia.

Usa Dataflow cuando haya streaming con semántica de tiempo del evento, transformaciones complejas o dependientes de estado, o cuando el equipo prefiera código versionable con pruebas automáticas.

Y hay una combinación muy razonable que se ve mucho en la práctica: Datastream para replicar el crudo + BigQuery SQL para transformar (ELT puro). Sin Data Fusion, sin Dataflow, sin instancias que pagar. Para muchos casos, incluido buena parte del de AlpinaShop, es la respuesta más eficiente.

  1. La decisión razonada de AlpinaShop

Con todo lo anterior sobre la mesa, Lucía y Marta deciden así:

El ERP de Sabadell → Datastream, no Data Fusion. El requisito es replicar tres tablas con baja latencia y capturar borrados. No hay transformación: los costes se calculan después con SQL en BigQuery, donde Lucía se maneja perfectamente. Datastream cuesta una fracción de una instancia Enterprise y no hay instancia que apagar. Las transformaciones del apartado 8 (mayúsculas, céntimos a euros, etiquetas de tipo de movimiento) se convierten en una vista de BigQuery, que además queda versionada en Git.

El CSV mensual del transportista → Data Fusion, con instancia efímera. Aquí sí gana: el fichero es sucio de verdad y las once directivas de Wrangler se escriben en veinte minutos viendo los datos, frente a un par de días de desarrollo y pruebas en Beam. Es mensual, así que la instancia se enciende, se ejecuta y se apaga. Con edición Basic y menos de 120 horas al mes, el coste es cero. La instancia se gestiona así:

# Apagar la instancia cuando no se usa (Enterprise y Basic lo permiten)
gcloud data-fusion instances update alpinashop-fusion \
  --location=europe-west1 --enable-instance-stopped

# Y volver a encenderla el dia de la carga mensual
gcloud data-fusion instances update alpinashop-fusion \
  --location=europe-west1 --no-enable-instance-stopped

La API del proveedor de mochilas → una Cloud Function programada. Son 400 productos en JSON, una vez a la semana. Levantar un clúster de Dataproc de tres nodos durante cinco minutos para leer 400 registros es desproporcionado. Cuarenta líneas de Python en Cloud Functions (06-03) disparadas por Cloud Scheduler resuelven el caso por céntimos.

La regla general que AlpinaShop adopta, y que sirve como criterio reutilizable:

  1. ¿El dato ya está en Cloud Storage con forma de tabla? → bq load.
  2. ¿Hay que replicar una base de datos sin transformar? → Datastream + vistas de BigQuery.
  3. ¿Es un volumen pequeño de una API? → Cloud Function programada.
  4. ¿Es un fichero sucio, recurrente, que va a limpiar un analista? → Data Fusion con instancia efímera.
  5. ¿Hay streaming, ventanas o lógica compleja? → Dataflow.
  6. ¿Es algorítmico con bibliotecas de Spark? → Dataproc Serverless.

Esa lista, con la disciplina de recorrerla en orden y quedarse en la primera que sirva, es lo que impide que una plataforma de datos de una pyme acabe costando lo que la de una multinacional.

Errores Comunes y Consejos

Dejar la instancia encendida. Es el error caro de este servicio y merece repetirse: una instancia Basic olvidada cuesta más de mil euros al mes sin ejecutar un solo pipeline. Apágala o bórrala.

Elegir Enterprise sin necesitarlo. Solo hace falta para CDC, alta concurrencia y alta disponibilidad. Basic cubre la mayoría de los casos de una pyme a menos de la mitad de precio.

Usar Data Fusion para volúmenes minúsculos. Un clúster de Dataproc para procesar 400 filas es desproporcionado. El arranque tarda más que el trabajo.

No versionar los pipelines. Aunque sean visuales, deben exportarse a JSON y guardarse en Git. Sin eso, una instancia borrada se lleva meses de trabajo y no hay historial de cambios.

# Exportar un pipeline para versionarlo
curl -H "Authorization: Bearer $(gcloud auth print-access-token)" \
  "${INSTANCIA}/v3/namespaces/default/apps/erp-costes-a-bigquery" \
  > pipelines/erp-costes-a-bigquery.json

Paralelizar la extracción sin medir el impacto en el origen. Number of Splits = 8 sobre el MySQL del ERP puede dejar sin servicio la báscula del almacén. Empieza por 1 y sube midiendo.

Escribir credenciales en la configuración del pipeline. Usa macros resueltas desde Secret Manager. Un pipeline exportado a JSON con la contraseña dentro acaba en Git, y eso es un incidente de seguridad.

Olvidar la paginación en el plugin HTTP. Se traen los primeros 100 registros y nadie lo nota hasta que faltan productos en un informe.

Saltarse el Preview. Es gratis, tarda segundos y muestra exactamente qué sale de cada bloque. Desplegar sin previsualizar es lanzar un clúster para descubrir un error de tipos.

No usar Upsert en el destino. Con Insert, cada reejecución duplica filas. Con Upsert sobre la clave de negocio, el pipeline es idempotente y se puede relanzar sin miedo.

Consejo: usa macros desde el principio. Fechas, rutas, nombres de tabla y credenciales como ${variable}. Convierte un pipeline rígido en uno reutilizable y orquestable.

Consejo: pon siempre un Error Collector. Los registros rechazados deben ir a algún sitio con su motivo. Es la misma disciplina que la cuarentena de Beam en 04-02.

Consejo: revisa el linaje después de cada despliegue. Es la mejor forma de verificar que el pipeline hace lo que crees que hace, y no cuesta nada.

Ejercicios

Ejercicio 1: directivas de Wrangler para el fichero del proveedor

El proveedor de mochilas envía este fichero de texto separado por tabuladores:

REF	DESCRIPCION	PVP_RECOMENDADO	COSTE	STOCK	ALTA	ACTIVO
mb-4001	Mochila Trekking 40L Azul	89,90 EUR	52,30 EUR	120	01-03-2024	S
MB-4002	  mochila trekking 30l roja  	74,50 EUR	43,10 EUR	0	15-06-2025	S
MB-4003	Mochila Alpina 55L	129,00 EUR	78,00 EUR	-	N/D	N

Escribe las directivas de Wrangler que produzcan un conjunto limpio con: sku en mayúsculas sin espacios; nombre con formato de título; pvp_eur y coste_eur como números decimales sin la palabra EUR; stock como entero, con - convertido en 0; fecha_alta como fecha yyyy-MM-dd, tolerando N/D como nulo; activo como booleano; y una columna calculada margen_pct. Los registros sin REF deben desviarse a error.

Ejercicio 2: diseñar el pipeline de envíos

Diseña el pipeline completo que lee el CSV mensual del transportista desde gs://alpinashop-datalake/transportista/2026/03/envios.csv, lo limpia, lo cruza con la tabla pedidos de alpinashop_analitica para añadir el país y el canal, calcula el retraso en días entre fecha_pedido y fecha_entrega, y escribe en alpinashop_analitica.envios. Indica: cada bloque con su tipo y configuración esencial, cómo tratar los envíos cuyo pedido_id no exista en pedidos, qué operación usar en el destino y por qué, y qué macros definirías.

Ejercicio 3: la decisión de herramienta, con justificación económica

AlpinaShop absorbe a un competidor y hereda cuatro integraciones nuevas. Para cada una, elige entre bq load, Datastream, Data Fusion, Dataflow, Dataproc Serverless o Cloud Function, y justifica incluyendo una estimación del orden de magnitud del coste mensual:

  1. Un PostgreSQL 14 de 80 GB con el histórico de clientes del competidor, que debe replicarse en BigQuery con menos de 5 minutos de retraso y capturando borrados.
  2. Un fichero Excel semanal de 3.000 filas que envía el equipo de compras, con columnas que cambian de nombre cada dos por tres y datos escritos a mano.
  3. Un flujo de 2.000 eventos por segundo de una aplicación móvil que hay que agregar por ventanas de 5 minutos antes de guardarlo.
  4. Un volcado nocturno en Parquet de 12 GB que el competidor ya deja en un bucket, con el esquema exacto de la tabla destino.

Soluciones

Solución 1

-- 1) Parsear el fichero separado por tabuladores, con cabecera
parse-as-csv :body '\t' true
drop :body

-- 2) SKU: quitar espacios y normalizar a mayusculas
trim :REF
uppercase :REF
rename :REF sku

-- 3) Desviar a error los registros sin referencia
send-to-error sku == null || sku.isEmpty()

-- 4) Nombre: limpiar espacios y aplicar formato de titulo
trim :DESCRIPCION
titlecase :DESCRIPCION
rename :DESCRIPCION nombre

-- 5) Importes: quitar ' EUR', cambiar la coma decimal y convertir a numero
find-and-replace :PVP_RECOMENDADO s/\s*EUR\s*//g
find-and-replace :COSTE s/\s*EUR\s*//g
find-and-replace :PVP_RECOMENDADO s/,/./g
find-and-replace :COSTE s/,/./g
set-type :PVP_RECOMENDADO double
set-type :COSTE double
rename :PVP_RECOMENDADO pvp_eur
rename :COSTE coste_eur

-- 6) Stock: '-' significa cero, no nulo
find-and-replace :STOCK s/^-$/0/g
set-type :STOCK int
rename :STOCK stock

-- 7) Fecha: 'N/D' a nulo, y formato dd-MM-yyyy a fecha real
find-and-replace :ALTA s/^N\/D$//g
parse-as-simple-date :ALTA dd-MM-yyyy
format-date :ALTA yyyy-MM-dd
rename :ALTA fecha_alta

-- 8) Activo: 'S'/'N' a booleano
set-column :ACTIVO (ACTIVO == 'S')
rename :ACTIVO activo

-- 9) Columna calculada: margen porcentual, protegido de la division por cero
set-column :margen_pct (pvp_eur != null && pvp_eur > 0) ? ((pvp_eur - coste_eur) / pvp_eur * 100) : null

Resultado esperado sobre las tres filas de ejemplo:

sku nombre pvp_eur coste_eur stock fecha_alta activo margen_pct
MB-4001 Mochila Trekking 40L Azul 89.90 52.30 120 2024-03-01 true 41.8
MB-4002 Mochila Trekking 30L Roja 74.50 43.10 0 2025-06-15 true 42.1
MB-4003 Mochila Alpina 55L 129.00 78.00 0 null false 39.5

Los tres puntos que se evalúan: distinguir - (que significa cero unidades) de N/D (que significa dato desconocido, es decir, nulo) —confundirlos falsearía cualquier informe de stock—; proteger la división de la columna calculada; y desviar a error en lugar de filtrar en silencio, para que un fichero con referencias vacías deje rastro.

Solución 2

Bloque 1 — Source: GCS

  • Path: gs://alpinashop-datalake/transportista/${anyo}/${mes}/envios.csv
  • Format: text (una fila por línea; el parseo se hace en Wrangler porque el separador es ;)
  • Service Account: sa-fusion-pipelines@...

Bloque 2 — Transform: Wrangler Las directivas del apartado 7, incluyendo el drop de nombre y apellidos por minimización de datos, y send-to-error para las filas sin pedido_id.

Bloque 3 — Source: BigQuery

  • Dataset/Table: alpinashop_analitica.pedidos
  • Import Query (mejor que leer la tabla entera):
SELECT pedido_id, fecha_pedido, envio.pais AS pais, canal
FROM `alpinashop-datos.alpinashop_analitica.pedidos`
WHERE fecha_pedido BETWEEN DATE '${fecha_desde}' AND DATE '${fecha_hasta}'

El filtro de partición es obligatorio: la tabla tiene require_partition_filter=TRUE desde 04-01, así que sin él el pipeline fallaría. Y aunque no lo tuviera, leer la tabla entera cada mes sería tirar dinero.

Bloque 4 — Analytics: Joiner

  • Inputs: salida de Wrangler (izquierda) y salida de BigQuery (derecha)
  • Join type: Left Outer sobre pedido_id
  • Fields: todos los del CSV limpio, más fecha_pedido, pais y canal del lado derecho

El tipo de unión es la decisión clave del ejercicio. Con Inner, los envíos cuyo pedido_id no exista en pedidos desaparecerían sin dejar rastro, y nadie sabría que faltan envíos. Con Left Outer se conservan todos, con pais y canal nulos, y esos nulos son la señal de que hay un problema de integridad que investigar: pedidos de meses fuera del rango de fechas, o identificadores mal escritos en el fichero del transportista.

Bloque 5 — Transform: Wrangler (segundo)

-- Retraso en dias entre pedido y entrega
set-column :retraso_dias (fecha_entrega != null && fecha_pedido != null) ? (dateDiff(fecha_entrega, fecha_pedido)) : null

-- Marca de envio huerfano, para poder contarlos en un informe de calidad
set-column :sin_pedido_asociado (pais == null)

Bloque 6 — Sink: BigQuery

  • Table: alpinashop_analitica.envios
  • Operation: Upsert con clave envio_id
  • Partition field: fecha_entrega; Cluster field: pais

Por qué Upsert. El fichero del transportista suele reenviarse corregido cuando hay errores, y a veces incluye envíos del mes anterior que se entregaron tarde. Con Insert, cada reenvío duplicaría filas y el informe de plazos mentiría. Con Upsert sobre envio_id, reejecutar el pipeline las veces que haga falta produce siempre el mismo resultado: es idempotente, exactamente el mismo principio que exigimos a los consumidores de Pub/Sub en 04-04.

Bloque 7 — Error Collector → Sink GCS

  • Path: gs://alpinashop-datalake/cuarentena/envios/${anyo}/${mes}/
  • Format: json

Macros a definir: ${anyo}, ${mes}, ${fecha_desde}, ${fecha_hasta}. Con ellas, el mismo pipeline sirve para cualquier mes y puede reprocesarse un histórico completo cambiando solo los argumentos. Sin ellas, habría que duplicar el pipeline o editarlo cada mes.

Solución 3

1. PostgreSQL de 80 GB replicado en BigQuery, menos de 5 minutos de retraso, con borrados → Datastream. Es literalmente su definición: CDC gestionado sobre PostgreSQL, sin código, con captura de DELETE —que una carga incremental por fecha nunca detectaría—. Data Fusion Enterprise también podría, pero exigiría la edición cara (~3.000 $/mes de instancia) para hacer exactamente lo mismo. Datastream se factura por GB procesado: la carga inicial de 80 GB más los cambios diarios sitúan el coste en el orden de unas decenas de euros al mes. La transformación posterior se hace con vistas en BigQuery, gratis en esfuerzo y versionables en Git.

2. Excel semanal de 3.000 filas con columnas cambiantes y datos a mano → Data Fusion con instancia efímera. Es el caso canónico de Wrangler: datos sucios escritos por personas, esquema inestable, y un analista —no un programador— que necesita corregirlo viendo los datos. Programarlo en Beam significaría redesplegar código cada vez que compras renombre una columna. Con instancia Basic encendida una hora a la semana, son 4 horas al mes, muy dentro de las 120 gratuitas: coste efectivo cero, más unos céntimos de Dataproc por ejecución. La clave de la respuesta es la palabra efímera: si la instancia se deja encendida, la misma solución cuesta más de mil euros al mes.

3. 2.000 eventos por segundo agregados en ventanas de 5 minutos → Dataflow. Streaming con ventanas temporales: el territorio exclusivo de Beam. Data Fusion no hace esto bien, bq load no aplica y una Cloud Function no puede mantener estado de ventana. Un pipeline de streaming con 2-3 trabajadores permanentes está en el orden de 100-150 € al mes, coste que hay que asumir conscientemente porque es la única herramienta que resuelve el requisito. Si el negocio tolerase 15 minutos de retraso en lugar de 5, la alternativa —suscripción de Pub/Sub a Cloud Storage más micro-lotes— costaría una fracción; merece la pena preguntarlo antes de encender el pipeline.

4. Parquet nocturno de 12 GB con el esquema exacto de la tabla destino → bq load. Sin transformación y con el esquema ya correcto, no hay nada que procesar. La carga por lotes en BigQuery es gratuita y Parquet lleva el esquema incorporado, así que ni siquiera hay que declararlo. Un comando en el orquestador de 04-06, disparado por la notificación del bucket que montamos en 04-04. Coste: 0 € de proceso, solo el almacenamiento de los 12 GB. Cualquier otra opción de esta lista sería pagar por hacer una copia, y es exactamente el reflejo que hay que corregir: la primera pregunta ante cualquier integración es "¿basta con bq load?".

Conclusión

Has visto el lado menos vistoso y más real de una plataforma de datos: los datos que no nacen en la nube. El MySQL del ERP en Sabadell, el CSV del transportista con sus fechas europeas y sus nulos escritos como guion, y la API del proveedor de mochilas. Ninguno publica en Pub/Sub, ninguno escribe Parquet, y Lucía necesita los tres para calcular el margen real por producto y explicar por qué caen las opiniones en ciertas zonas.

Has entendido qué es un ETL/ELT visual y, sobre todo, a quién sirve: a perfiles que dominan el negocio y el SQL pero no van a escribir Beam, y a organizaciones que necesitan trazabilidad documentada. Y también a quién no sirve, que es igual de importante. Sabes que Data Fusion es CDAP gestionado, con su Studio, su Wrangler, su Hub y su registro de linaje, y que por debajo no ejecuta nada: traduce el grafo visual a Spark y levanta un Dataproc efímero, lo que explica sus 2-5 minutos de arranque y por qué no sirve para latencias bajas.

Has puesto el coste por delante en lugar de esconderlo: ediciones Developer, Basic y Enterprise, facturadas por hora de instancia existente, se use o no, con una instancia Basic permanente costando más que todo el resto de la plataforma de AlpinaShop junta. Ese dato no descalifica el producto —para una empresa con veinte orígenes es barato frente a los salarios que ahorra— pero sí obliga a la estrategia de instancia efímera en una pyme.

Has limpiado el CSV del transportista con Wrangler y sus directivas: parseo por punto y coma, trim, nulos escritos como guion, fechas dd/MM/yyyy convertidas a tipo fecha, decimales con coma, división de la columna de destinatario, formato de título, columnas derivadas, y el drop final del nombre y apellidos por minimización de datos, con la advertencia expresa de RGPD y de revisión por compliance. Has valorado la ventaja real de la herramienta, que no es la potencia sino el ciclo de retroalimentación: ves el efecto de cada directiva sobre datos reales al instante.

Has diseñado el pipeline del ERP con sus macros de fecha, su lectura particionada con la prudencia de no tumbar la báscula del almacén, su limpieza, su agregación y su destino en Upsert para que sea idempotente y reejecutable. Conoces el Hub y sus conectores, el plugin HTTP con su paginación fácil de olvidar, el flujo Preview → Deploy → Run → Schedule, los perfiles de cómputo que controlan el clúster subyacente, y la API REST con la que el orquestador lo invocará. Has visto el linaje a nivel de campo, que responde a "¿de dónde sale esta columna?" y "¿qué se rompe si cambio esto?" sin días de arqueología, y la replicación CDC desde MySQL con su binlog_format = ROW innegociable y su gran argumento: es la única forma de enterarse de los borrados.

Y has cerrado con la tabla honesta y con la decisión razonada de AlpinaShop, que no es "usemos Data Fusion para todo" sino una lista ordenada: primero bq load si basta, luego Datastream si es replicar, luego una Cloud Function si es poca cosa, luego Data Fusion con instancia efímera si el fichero es sucio y lo limpia un analista, y Dataflow o Dataproc si hay streaming o algoritmos. Recorrer esa lista en orden y quedarse en la primera opción que sirva es lo que separa una plataforma de datos proporcionada de una carísima.

Y ahora aparece un problema nuevo, que es consecuencia directa del éxito de las cuatro lecciones anteriores. AlpinaShop tiene una exportación nocturna de Cloud SQL, una carga en BigQuery, un pipeline de Dataflow, un trabajo de Spark en Dataproc Serverless, un pipeline mensual de Data Fusion, varias consultas de agregación y una vista materializada que refrescar. Son siete u ocho procesos que dependen unos de otros: no tiene sentido lanzar la agregación antes de que haya terminado la carga, ni refrescar la vista de negocio con datos a medias.

Ahora mismo, eso lo resuelve un cron en una máquina virtual que lanza scripts a horas fijas, calculadas a ojo con margen de sobra. Funciona hasta la primera noche en que la exportación tarda veinte minutos más de lo normal: entonces la carga lee un fichero incompleto, la agregación calcula sobre datos parciales, el informe de dirección amanece con cifras falsas, y nadie se entera hasta que alguien las mira a media mañana.

En 04-06, Cloud Composer y Workflows, resolveremos eso. Verás por qué un cron no es un orquestador y qué significa realmente coordinar dependencias, reintentos, alertas y reprocesos. Conocerás Cloud Composer —Apache Airflow gestionado— con sus DAG, tareas, operadores y sensores, y escribirás el DAG completo del proceso nocturno de AlpinaShop comentado línea a línea, con la advertencia clara de que Composer es caro para una pyme. Y conocerás Workflows, la orquestación serverless en YAML sin coste fijo, con Cloud Scheduler para el disparo, que es por donde AlpinaShop va a empezar.

Curso de Google Cloud Platform (GCP)

Módulo 1: Introducción a Google Cloud Platform

Módulo 2: Servicios principales de GCP

Módulo 3: Redes y seguridad

Módulo 4: Datos y análisis

Módulo 5: Aprendizaje automático e IA

Módulo 6: DevOps y monitoreo

Módulo 7: Temas avanzados de GCP

Módulo 8: Proyecto final

© Copyright 2026. Todos los derechos reservados