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 columnaDESTINATARIOque 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
- El problema de los datos que no nacen en Google Cloud
- Qué es un ETL/ELT visual y a quién sirve
- Qué es Data Fusion: CDAP gestionado
- Ediciones, coste por hora y sus consecuencias
- Crear la instancia y entender qué se ha creado
- El Studio: orígenes, transformaciones y destinos
- Wrangler: limpiar el CSV del transportista viendo los datos
- El pipeline real: MySQL del ERP a
alpinashop_analitica - Conectores y plugins del Hub
- Desplegar, ejecutar y programar
- Qué ocurre por debajo: Dataproc efímero
- Linaje de datos a nivel de campo
- Replicación con CDC desde el MySQL del ERP
- Cuándo Data Fusion, cuándo Dataflow, cuándo Datastream, cuándo
bq load - La decisión razonada de AlpinaShop
- 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.
- 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.
- 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.
- 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.
- 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=analiticaLa 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.
- 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 Collectorrecoge los registros que unWranglerno 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.
- 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:
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.
- El pipeline real: MySQL del ERP a
alpinashop_analitica
alpinashop_analiticaVamos 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 $CONDITIONSDos 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.$CONDITIONSconnumSplits: si configurasSplit-By Field Name = skuyNumber 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_compradasAvg(coste_unitario_eur)→coste_medio_eurMin(fecha_movimiento)→primera_compraMax(fecha_movimiento)→ultima_compraCount(*)→num_movimientos
Bloque 4 — Sink: BigQuery.
- Dataset:
alpinashop_analitica - Table:
costes_producto_erp - Operation:
Upsertcon clavesku(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 accionableEsa 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.
- 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 HeaderoIncrement 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.
- 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:
- 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
DirectRunnerde Beam y hay que usarlo siempre antes de desplegar. - Deploy. Congela el pipeline con un número de versión. Para cambiarlo, se crea una versión nueva.
- Run. Ejecuta. Se pueden pasar argumentos en tiempo de ejecución (las macros).
- 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.
- 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.
- 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.
- 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):
- Origen: MySQL, con la conexión y el usuario
cdc_datafusion. - Selección de tablas:
movimientos_stock,proveedores,articulos. - Destino: BigQuery, dataset
alpinashop_analitica, con un prefijo de staging. - Evaluación de la fuente: comprueba permisos y configuración antes de empezar.
- 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.
- Cuándo Data Fusion, cuándo Dataflow, cuándo Datastream, cuándo
bq load
bq loadLa 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.
- 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-stoppedLa 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:
- ¿El dato ya está en Cloud Storage con forma de tabla? →
bq load. - ¿Hay que replicar una base de datos sin transformar? → Datastream + vistas de BigQuery.
- ¿Es un volumen pequeño de una API? → Cloud Function programada.
- ¿Es un fichero sucio, recurrente, que va a limpiar un analista? → Data Fusion con instancia efímera.
- ¿Hay streaming, ventanas o lógica compleja? → Dataflow.
- ¿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.jsonParalelizar 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:
- 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.
- 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.
- 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.
- 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 Outersobrepedido_id - Fields: todos los del CSV limpio, más
fecha_pedido,paisycanaldel 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:
Upsertcon claveenvio_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
- ¿Qué es Google Cloud Platform?
- Configuración de tu cuenta de GCP
- Descripción general de la consola de GCP
- Proyectos, jerarquía de recursos y facturación
- Regiones, zonas y modelo de responsabilidad compartida
- Cloud Shell y la CLI de gcloud
Módulo 2: Servicios principales de GCP
- Compute Engine: máquinas virtuales en Google Cloud
- Cloud Storage: almacenamiento de objetos
- Cloud SQL: bases de datos relacionales gestionadas
- App Engine: plataforma como servicio
- Google Kubernetes Engine (GKE)
- Bases de datos NoSQL: Firestore, Bigtable y Spanner
- Cómo elegir el servicio de cómputo adecuado
Módulo 3: Redes y seguridad
- Redes VPC
- Balanceo de carga en la nube
- Cloud CDN
- Gestión de identidad y acceso (IAM)
- Cloud Armor
- Secretos y cifrado: Secret Manager y Cloud KMS
- Cloud DNS, certificados TLS y publicación segura de servicios
Módulo 4: Datos y análisis
- BigQuery: el almacén de datos analítico
- Cloud Dataflow: procesamiento de datos por lotes y en streaming
- Cloud Dataproc: Spark y Hadoop gestionados
- Cloud Pub/Sub: mensajería asíncrona
- Cloud Data Fusion: integración de datos sin código
- Orquestación de pipelines con Cloud Composer y Workflows
- Gobierno del dato y cuadros de mando con Dataplex y Looker Studio
Módulo 5: Aprendizaje automático e IA
- Vertex AI: la plataforma de machine learning de GCP
- AutoML: modelos a medida sin escribir código
- TensorFlow en GCP: entrenamiento y servicio de modelos
- API de lenguaje natural
- API de visión
- IA generativa en Vertex AI: modelos Gemini y embeddings
- MLOps: del modelo al producto con Vertex AI Pipelines
Módulo 6: DevOps y monitoreo
- Cloud Build: integración continua en GCP
- Cloud Source Repositories y gestión del código fuente
- Cloud Functions: funciones sin servidor
- Cloud Monitoring (antes Stackdriver): métricas, paneles y alertas
- Cloud Deployment Manager e infraestructura como código nativa
- Cloud Logging y Cloud Trace: logs, trazas y diagnóstico
- Terraform en GCP: infraestructura como código en la práctica
Módulo 7: Temas avanzados de GCP
- Híbrido y multinube con Anthos
- Computación sin servidor con Cloud Run
- Redes avanzadas: VPC compartida, peering y conectividad híbrida
- Mejores prácticas de seguridad
- Gestión y optimización de costos
- Fiabilidad: SLO, alta disponibilidad y recuperación ante desastres
- Gobierno a escala: organización, políticas y auditoría
