Tot el que hem construït fins ara comparteix una suposició còmoda: que les dades ja són a Google Cloud. L'històric venia de Cloud SQL, els esdeveniments arriben per Pub/Sub, les exportacions aterren a alpinashop-catalogo. Però AlpinaShop, com qualsevol empresa real, té dades que viuen en altres llocs i arriben d'altres maneres.
Són tres casos concrets, i cap no és exòtic:
- L'ERP del magatzem és un MySQL 8 que corre en un servidor físic a les oficines de Sabadell. Conté l'estoc real, les entrades de mercaderia, els proveïdors i els costos de compra. Ningú no el migrarà aquest any: funciona, està integrat amb la bàscula i amb la impressora d'etiquetes, i el proveïdor cobra per cada canvi.
- El transportista envia cada mes un CSV per correu electrònic amb els enviaments, les incidències i els terminis de lliurament reals. Separat per punt i coma, amb dates en format
dd/mm/aaaa, imports amb coma decimal, noms de ciutat en majúscules i sense accents, i una columnaDESTINATARIOque barreja nom i cognoms. - El proveïdor de motxilles ofereix una API REST amb el seu catàleg actualitzat: referències, preus de cost, disponibilitat i fitxes tècniques.
La Lucía necessita els tres. Necessita creuar els costos de compra de l'ERP amb les vendes de BigQuery per saber el marge real per producte. Necessita els terminis del transportista per explicar per què cauen les ressenyes en certes zones. I necessita el catàleg del proveïdor per detectar productes descatalogats que continuen publicats al web.
La Lucía sap SQL. Sap molt SQL. Però no ha escrit un pipeline d'Apache Beam a la vida i no començarà ara, i el Dani té la seva pròpia llista de tasques. Si cada fitxer nou requereix un desenvolupament, aquestes dades no arribaran mai.
Cloud Data Fusion existeix per a aquest buit: una eina d'integració de dades visual, on es construeixen pipelines arrossegant blocs i es netegen dades veient-les, sense escriure codi. En aquesta lliçó la faràs servir per construir el pipeline de l'ERP, netejaràs el CSV del transportista amb Wrangler, entendràs el llinatge a nivell de camp que és la seva virtut més gran, muntaràs la replicació amb captura de canvis des de MySQL i —això és igual d'important— aprendràs quan no utilitzar-la, perquè té un cost per hora que condiciona tota la decisió.
Contingut
- El problema de les dades que no neixen a Google Cloud
- Què és un ETL/ELT visual i a qui serveix
- Què és Data Fusion: CDAP gestionat
- Edicions, cost per hora i les seves conseqüències
- Crear la instància i entendre què s'ha creat
- L'Studio: orígens, transformacions i destinacions
- Wrangler: netejar el CSV del transportista veient les dades
- El pipeline real: MySQL de l'ERP a
alpinashop_analitica - Connectors i plugins del Hub
- Desplegar, executar i programar
- Què passa per sota: Dataproc efímer
- Llinatge de dades a nivell de camp
- Replicació amb CDC des del MySQL de l'ERP
- Quan Data Fusion, quan Dataflow, quan Datastream, quan
bq load - La decisió raonada d'AlpinaShop
- El problema de les dades que no neixen a Google Cloud
Se'n diu integració de dades, i és la feina menys glamurosa i que més temps consumeix de qualsevol projecte analític. Les enquestes del sector ho repeteixen sense variació: entre el 60 % i el 80 % de l'esforç d'un projecte de dades se'n va a aconseguir que les dades arribin, no a analitzar-les.
Els problemes concrets, amb els exemples d'AlpinaShop:
| Problema | Cas a AlpinaShop |
|---|---|
| Connectivitat | El MySQL és a Sabadell, darrere d'un tallafocs, sense IP pública |
| Formats | CSV amb ;, dates dd/mm/aaaa, decimals amb coma, sense capçalera fiable |
| Qualitat | Ciutats en majúscules sense accents, nuls escrits com a -, espais de sobres |
| Esquemes canviants | El transportista va afegir una columna al gener sense avisar |
| Freqüència dispar | L'ERP canvia contínuament; el CSV arriba mensualment; l'API es consulta a demanda |
| Volum | L'històric de l'ERP són 12 milions de moviments d'estoc |
| Traçabilitat | Ningú no sap d'on va sortir el camp coste_unitario d'un informe del 2025 |
Cadascun d'aquests problemes es pot resoldre programant. La qüestió és que resoldre'ls programant per a cada origen, i mantenir aquest codi quan el transportista canvia el format, és una feina contínua que una pime de 40 persones no es pot permetre dedicar al Dani.
- Què és un ETL/ELT visual i a qui serveix
ETL significa extreure, transformar, carregar: es treu la dada de l'origen, es transforma fora, i es carrega ja neta a la destinació. ELT inverteix els dos últims: es carrega la dada en cru i es transforma dins de la destinació, aprofitant la seva potència.
| Enfocament | On es transforma | Avantatge | Inconvenient |
|---|---|---|---|
| ETL | En un motor intermedi | La destinació només rep dades netes | Aquest motor s'ha de dimensionar i pagar |
| ELT | A la destinació (BigQuery) | Aprofita el seu motor; conserves el cru | La destinació guarda dades brutes; cost de consulta |
Amb magatzems moderns com BigQuery, la tendència clara és ELT: carregar en cru i transformar amb SQL. Però hi ha una part que continua sent ETL inevitablement: treure la dada de l'origen i portar-la fins a la porta. Això és l'extracció, i és exactament on Data Fusion aporta.
Un ETL visual és una eina on el pipeline es construeix amb un llenç i blocs en comptes de codi. A qui serveix?
Serveix molt bé a:
- Analistes com la Lucía, que coneixen el negoci i les dades però no programen pipelines distribuïts.
- Equips petits sense enginyers de dades dedicats.
- Integracions estàndard: llegir una taula, netejar uns camps, escriure en una altra taula.
- Organitzacions que necessiten traçabilitat documentada d'on surt cada dada (auditories, compliment normatiu).
Serveix malament a:
- Lògica de negoci complexa amb condicions imbricades i estat.
- Equips que ja tenen enginyers i control de versions madur: un pipeline visual es versiona pitjor que un fitxer
.py. - Streaming amb finestres i temps de l'esdeveniment: per a això hi ha Beam.
- Pressupostos ajustats amb ús esporàdic, pel motiu de l'apartat 4.
I un advertiment honest que convé fer des del principi: sense codi no vol dir sense coneixement. Per construir un pipeline a Data Fusion cal entendre esquemes, tipus, unions, claus i particions exactament igual que programant-lo. El que s'estalvia és la sintaxi i la infraestructura, no el pensament.
- Què és Data Fusion: CDAP gestionat
Cloud Data Fusion és la versió gestionada de CDAP (Cask Data Application Platform), una plataforma open source d'integració de dades que Google va adquirir i ofereix com a servei.
Els seus components:
| Component | Què fa |
|---|---|
| Studio | El llenç visual on es dibuixa el pipeline |
| Wrangler | Explorador i netejador interactiu de dades |
| Hub | Catàleg de plugins, connectors i pipelines d'exemple |
| Metadades i llinatge | Registre automàtic de quin camp ve d'on |
| Replicació | Mòdul de CDC per copiar bases de dades en continu |
| Motor d'execució | Genera i executa la feina a Dataproc |
L'última fila és la més important per entendre el producte: Data Fusion no executa res per si mateix. Tradueix el pipeline visual a una feina de Spark o MapReduce i l'executa en un clúster de Dataproc que crea al vol. Això explica el seu comportament, el seu temps d'arrencada i bona part del seu cost, i ho desenvoluparem a l'apartat 11.
Ser CDAP gestionat té una altra conseqüència rellevant: els pipelines són portables. Un pipeline exportat de Data Fusion es pot importar en un CDAP autogestionat a qualsevol lloc. No és un format propietari tancat.
- Edicions, cost per hora i les seves conseqüències
Aquí hi ha la característica que condiciona totes les decisions sobre aquest producte, i cal posar-la al davant en comptes d'amagar-la al final.
| Edició | Per a què | Cost aproximat per hora d'instància | Notes |
|---|---|---|---|
| Developer | Proves i desenvolupament | ~0,35 $ | Sense alta disponibilitat, capacitat limitada |
| Basic | Producció senzilla | ~1,80 $ | Inclou 120 hores gratuïtes al mes per compte |
| Enterprise | Producció exigent | ~4,20 $ | Alta disponibilitat, més concurrència, llinatge complet, CDC |
Verifica els preus vigents a la documentació oficial; el que importa és el model, no la xifra exacta.
I el model és aquest: es paga per hora d'instància existent, s'utilitzi o no. No per pipeline executat, ni per dada processada. La instància és un entorn que està encès.
Fem el càlcul per a AlpinaShop, perquè és la conversa que cal tenir amb direcció:
| Escenari | Hores/mes | Cost aproximat |
|---|---|---|
| Instància Basic permanent | 730 | ~1.310 $ (menys 120 h gratis: ~1.100 $) |
| Instància Enterprise permanent | 730 | ~3.070 $ |
| Instància Developer permanent | 730 | ~255 $ |
| Basic encesa 4 h al dia | 120 | 0 $ (dins de les 120 h gratuïtes) |
I a això cal sumar-hi el cost del clúster de Dataproc que s'aixeca per executar cada pipeline, que a la pràctica sol ser menor però no és zero.
La conseqüència és directa i cal dir-la sense adorns: una instància de Data Fusion permanent costa més al mes que tota la resta de la plataforma de dades d'AlpinaShop junta. BigQuery costa uns pocs euros, Pub/Sub cèntims, Dataflow per lots cèntims. Data Fusion Basic permanent costaria més de mil euros.
Això no invalida el producte: per a una empresa amb vint orígens de dades i un equip d'integració, mil euros al mes és una ganga davant dels sous que estalvia. Però per a una pime amb tres orígens, l'equació no surt, i cal dir-ho.
L'estratègia que sí que funciona en una pime és tractar la instància com els clústers de Dataproc de 04-03: efímera. S'encén per desenvolupar, s'apaga en acabar; i els pipelines ja desplegats s'executen igual, perquè la feina la fa Dataproc. Hi tornarem a l'apartat 15.
- Crear la instància i entendre què s'ha creat
gcloud config set project alpinashop-datos
gcloud services enable datafusion.googleapis.com
# Instancia de desenvolupament: la mes barata per aprendre i 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ó triga entre 15 i 25 minuts. No és un error: s'està aprovisionant un entorn complet de CDAP en un projecte gestionat per Google.
Aquest detall del "projecte gestionat" importa per als permisos. Data Fusion crea la instància en un projecte propi de Google (el tenant project) i actua sobre el teu mitjançant un agent de servei, al qual cal concedir permisos explícitament:
PROY_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_FUSION="service-${PROY_NUM}@gcp-sa-datafusion.iam.gserviceaccount.com"
# L'agent necessita poder actuar sobre el projecte
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:${SA_FUSION}" \
--role="roles/datafusion.serviceAgent"
# Compte de servei que faran servir els clusters de Dataproc que executen els pipelines
gcloud iam service-accounts create sa-fusion-pipelines \
--display-name="Execucio 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
# L'agent de Data Fusion ha de poder fer servir aquest compte
gcloud iam service-accounts add-iam-policy-binding "$SA_PIPE" \
--member="serviceAccount:${SA_FUSION}" \
--role="roles/iam.serviceAccountUser"Accés a la interfície:
gcloud data-fusion instances describe alpinashop-fusion \
--location=europe-west1 --format="value(apiEndpoint, serviceEndpoint, state)"I la connectivitat amb Sabadell, que és el requisit real per llegir de l'ERP. Hi ha tres opcions, en ordre de preferència:
| Opció | Com funciona | Valoració |
|---|---|---|
| Cloud VPN o Interconnect | Túnel entre alpinashop-vpc i la xarxa de les oficines |
La correcta. El MySQL no s'exposa mai a internet |
| IP pública amb llista blanca i TLS | S'obre el port només a rangs concrets | Acceptable de mala gana; superfície d'atac innecessària |
| Exportar a fitxer i pujar-lo | Un script de l'ERP bolca CSV a Cloud Storage | Senzill, però no permet CDC ni dades fresques |
Per a AlpinaShop, la Marta munta un túnel de Cloud VPN amb encaminament cap a la subxarxa de l'ERP, coherent amb tot el que hem vist a 03-01: el MySQL continua sense IP pública i el trànsit va xifrat. La configuració detallada de la connectivitat híbrida és territori de 07-03.
- L'Studio: orígens, transformacions i destinacions
L'Studio és un llenç. A l'esquerra hi ha una paleta de blocs agrupats per categoria; s'arrosseguen al llenç i es connecten amb fletxes que representen el flux de les dades.
Els blocs es classifiquen així:
| Categoria | Què fa | Exemples |
|---|---|---|
| Source | Llegeix dades | Database (JDBC), BigQuery, GCS, Salesforce, HTTP, Kafka |
| Transform | Modifica registres | Wrangler, JavaScript, Python, Projection, Encoder |
| Analytics | Agrega i uneix | Group By, Joiner, Deduplicate, Distinct, Row Denormalizer |
| Conditions and Actions | Control de flux i accions | Condition, Email, BigQuery Execute, Database Execute |
| Sink | Escriu dades | BigQuery, GCS, Database, Spanner, Pub/Sub |
| Error Handlers | Recullen registres rebutjats | Error Collector |
Cada bloc té un panell de configuració: cadena de connexió, taula, esquema de sortida, opcions específiques. I cada bloc declara un esquema de sortida —la llista de camps amb els seus tipus— que es propaga al següent. Aquesta propagació és el que permet el llinatge de l'apartat 12.
El pipeline que construirem, dibuixat:
flowchart LR
S1["Source: Database<br/>MySQL ERP Sabadell<br/>taula movimientos_stock"]
S2["Source: GCS<br/>CSV del transportista"]
W1["Transform: Wrangler<br/>neteja i tipus"]
W2["Transform: Wrangler<br/>dates, decimals, noms"]
J["Analytics: Joiner<br/>per sku"]
G["Analytics: Group By<br/>cost mitja per sku"]
K1["Sink: BigQuery<br/>alpinashop_analitica.stock_erp"]
K2["Sink: BigQuery<br/>alpinashop_analitica.envios"]
E["Error Collector<br/>-> GCS quarantena"]
S1 --> W1 --> G --> K1
S2 --> W2 --> J
W1 --> J
J --> K2
W2 -.rebutjos.-> E
Fixa't en dues coses del diagrama, perquè reprodueixen les bones pràctiques de les lliçons anteriors:
- La branca de rebutjos existeix també aquí. El bloc
Error Collectorrecull els registres que unWranglerno va poder processar i els envia a una destinació de quarantena, exactament igual que les sortides etiquetades de Beam a 04-02. - Un mateix origen alimenta dues branques. El llenç és un graf dirigit, no una línia, igual que el graf de Beam.
- Wrangler: netejar el CSV del transportista veient les dades
Wrangler és la peça que justifica per si sola aprendre Data Fusion. És un entorn interactiu on carregues una mostra de les dades, la veus en forma de taula, i apliques transformacions —anomenades directives— veient-ne l'efecte immediatament.
Partim del CSV real del transportista, tal com arriba:
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
Els problemes salten a la vista: separador ;, dates europees, decimals amb coma, nuls com a -, espais sobrants, majúscules inconsistents, nom i cognoms junts, i un lliurament sense data.
Les directives de Wrangler s'escriuen una per línia i s'apliquen en ordre. Es poden generar des del menú contextual de cada columna, però és molt més ràpid escriure-les:
-- 1) Partir la linia pel separador i anomenar les columnes parse-as-csv :body ';' true drop :body -- 2) Netejar espais sobrants a totes les columnes de text trim :DESTINATARIO trim :CIUDAD trim :INCIDENCIA -- 3) Nuls: el transportista escriu '-' on no hi ha dada find-and-replace :FECHA_ENTREGA s/^-$//g find-and-replace :INCIDENCIA s/^-$//g set-column :INCIDENCIA (INCIDENCIA == null || INCIDENCIA.isEmpty()) ? null : INCIDENCIA -- 4) Dates europees -> tipus data real parse-as-simple-date :FECHA_ENTREGA dd/MM/yyyy format-date :FECHA_ENTREGA yyyy-MM-dd -- 5) Decimals amb coma -> punt, i despres a numero find-and-replace :PESO s/,/./g find-and-replace :IMPORTE s/,/./g set-type :PESO double set-type :IMPORTE double -- 6) Partir 'COGNOMS, NOM' en dues columnes split-to-columns :DESTINATARIO , rename :DESTINATARIO_1 apellidos rename :DESTINATARIO_2 nombre trim :apellidos trim :nombre -- 7) Normalitzar el format dels noms propis titlecase :apellidos titlecase :nombre uppercase :CIUDAD -- 8) Columna derivada: hi va haver incidencia si el camp no esta buit set-column :hubo_incidencia (INCIDENCIA != null) -- 9) Noms de columna en minuscules i coherents amb la resta del magatzem 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 files sense identificador de comanda: no serveixen de res filter-rows-on condition-false pedido_id != null && !pedido_id.isEmpty() -- 11) Treure la columna de nom i cognoms per MINIMITZACIO DE DADES drop :apellidos drop :nombre
Les directives més útils, agrupades:
| Família | Directives | Per a què |
|---|---|---|
| Parseig | parse-as-csv, parse-as-json, parse-as-xml, parse-as-fixed-length |
Convertir text cru en columnes |
| Dates | parse-as-simple-date, format-date, parse-as-datetime |
Formats regionals |
| Text | trim, uppercase, lowercase, titlecase, cleanse-column-names |
Normalització |
| Cerca | find-and-replace, extract-regex-groups, split-to-columns |
Expressions regulars |
| Tipus | set-type, fill-null-or-empty |
Conversió i nuls |
| Filtres | filter-rows-on, filter-row-if-matched |
Descartar files |
| Columnes | rename, drop, keep, set-column, merge |
Estructura |
| Qualitat | send-to-error |
Desviar registres invàlids |
Els passos 10 i 11 de l'exemple mereixen comentari.
El pas 10 utilitza filter-rows-on, que descarta silenciosament. Si prefereixes conservar el que es descarta per revisar-ho —i normalment hauries de fer-ho—, s'utilitza send-to-error, que envia el registre al Error Collector del pipeline en comptes de llençar-lo:
El pas 11 és una decisió de compliment normatiu, no de neteja:
Avís de RGPD. El CSV del transportista conté el nom i els cognoms del destinatari, que és una dada personal identificativa. Per a l'anàlisi de terminis de lliurament i incidències, aquesta dada no aporta absolutament res: n'hi ha prou amb el codi postal i l'identificador de comanda. Eliminar-la en la fase de neteja, abans que arribi al magatzem analític, és minimització de dades per disseny i és la pràctica correcta. Qualsevol tractament de dades personals reals ha de ser revisat per un professional de compliance o el DPD abans de passar a producció. Totes les dades d'aquest curs són fictícies.
L'avantatge de Wrangler davant d'escriure això en Python o SQL no és la potència —Python pot fer tot això i més—, sinó el cicle de retroalimentació: apliques una directiva i en veus immediatament l'efecte sobre 100 files reals. Quan parse-as-simple-date falla perquè hi ha una data amb format diferent a la fila 47, ho veus a l'instant, no quan el pipeline porti vint minuts en producció.
- El pipeline real: MySQL de l'ERP a
alpinashop_analitica
alpinashop_analiticaAnem amb la integració que de veritat vol la Lucía: portar els costos de compra i l'estoc de l'ERP per poder calcular el marge real per producte.
Bloc 1 — Source: Database. Es configura amb:
- Plugin type:
Database(genèric JDBC). - JDBC driver: el connector de MySQL, que cal pujar prèviament des del Hub o com a plugin propi.
- Connection string:
jdbc:mysql://10.20.0.15:3306/erp_almacen(IP privada abastable per la VPN). - Import Query: la consulta que extreu les dades.
-- Consulta d'extraccio de l'ERP.
-- $CONDITIONS es OBLIGATORI si s'utilitza lectura particionada: Data Fusion
-- el substitueix per un rang diferent a cada tasca paral·lela.
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 detalls fonamentals d'aquesta configuració:
${fecha_desde}i${fecha_hasta}són macros de Data Fusion: arguments en temps d'execució. Permeten que el mateix pipeline serveixi per a la càrrega inicial completa i per a la incremental diària, sense duplicar-lo. L'orquestrador de 04-06 els passarà els valors.$CONDITIONSambnumSplits: si configuresSplit-By Field Name = skuiNumber of Splits = 4, Data Fusion llança quatre consultes en paral·lel sobre rangs diferents. Això accelera moltíssim la càrrega inicial, però posa quatre vegades més càrrega sobre el MySQL de l'ERP. Amb una base de dades de producció que a més atén la bàscula del magatzem, cal ser prudent: un o dos splits, i executar de matinada.
Bloc 2 — Transform: Wrangler. Neteja específica de l'ERP:
-- L'ERP barreja majuscules i minuscules als SKU uppercase :sku trim :sku -- Codis de moviment interns -> etiquetes llegibles set-column :tipo_movimiento (tipo_movimiento == 'E' ? 'entrada' : (tipo_movimiento == 'S' ? 'salida' : 'ajuste')) -- L'ERP guarda els costos en centims com a enter set-column :coste_unitario_eur coste_unitario / 100.0 drop :coste_unitario -- Descartar moviments sense SKU: son ajustos comptables, no d'estoc send-to-error sku == null || sku.isEmpty() -- Marca de quan es va extreure, per poder auditar carregues set-column :cargado_en datetime:CurrentDateTime()
Bloc 3 — Analytics: Group By. Cost mitjà ponderat per 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
Bloc 4 — Sink: BigQuery.
- Dataset:
alpinashop_analitica - Table:
costes_producto_erp - Operation:
Upsertamb clausku(actualitza si existeix, insereix si no) - Truncate Table: desactivat
- Service Account:
[email protected] - Location:
europe-west1
L'opció Upsert és la que fa el pipeline reexecutable sense duplicar. És l'equivalent al MERGE de SQL, i és el que permet llançar el pipeline dues vegades sense espatllar res, la mateixa propietat d'idempotència que perseguíem a 04-04.
Bloc 5 — Error Collector → Sink GCS. Els registres desviats amb send-to-error van a gs://alpinashop-datalake/cuarentena/erp/, en format JSON, amb el motiu del rebuig.
I amb aquestes dades ja a BigQuery, la Lucía pot per fi respondre a la seva pregunta amb SQL:
-- Marge real per producte: preu de venda mitja contra cost de l'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; -- els pitjors marges primer: aixo es el que es accionableAquesta consulta, que ordena pel marge més baix, és exactament el tipus de resultat que canvia decisions: productes que es venen molt i deixen poc. No era possible abans d'aquesta lliçó perquè el cost vivia en un servidor de Sabadell.
- Connectors i plugins del Hub
El Hub és el catàleg des del qual s'instal·len plugins a la instància. Categories principals:
| Tipus | Exemples disponibles |
|---|---|
| Bases de dades | 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 |
| Fitxers i protocols | HTTP, FTP/SFTP, Amazon S3, Azure Blob, Excel, XML |
| Transformacions | Wrangler, JavaScript, Python Evaluator, XML Parser, Validator |
| Analytics | Joiner, Group By, Deduplicate, Pivot, Window Aggregation |
Per al tercer origen d'AlpinaShop, l'API del proveïdor de motxilles, s'utilitza el plugin HTTP:
- URL:
https://api.proveedor-montana.example/v2/catalogo - HTTP Method:
GET - Headers:
Authorization: Bearer ${api_token}— on${api_token}és una macro que es resol des de Secret Manager (03-06), mai escrita al pipeline. - Format:
json - JSON/XML Result Path:
$.productos— la ruta dins de la resposta on hi ha l'array. - Pagination Type:
Link HeaderoIncrement an Index, segons el que admeti l'API.
La paginació és el que més s'oblida: sense configurar-la, es porten només els primers N resultats i ningú no se n'adona fins que falten productes.
Si un origen no té plugin, queden tres sortides: escriure'l en Java (CDAP és extensible), utilitzar el Python Evaluator per a transformacions puntuals, o —el més raonable— exportar la dada a Cloud Storage amb un script i llegir-la des d'allà. L'última sol ser la resposta correcta per a un origen exòtic i de volum baix.
- Desplegar, executar i programar
Un pipeline a Data Fusion té dos estats: esborrany (s'edita, es previsualitza) i desplegat (immutable, executable, versionat).
El flux de treball:
- Preview. Executa el pipeline amb una mostra petita sense escriure a la destinació. Mostra les dades que surten de cada bloc. És l'equivalent al
DirectRunnerde Beam i cal utilitzar-lo sempre abans de desplegar. - Deploy. Congela el pipeline amb un número de versió. Per canviar-lo, es crea una versió nova.
- Run. Executa. Es poden passar arguments en temps d'execució (les macros).
- Schedule. Programa execucions periòdiques amb expressió cron.
Tot això també es fa per API REST, que és com ho invocarà l'orquestrador de 04-06:
INSTANCIA=$(gcloud data-fusion instances describe alpinashop-fusion \
--location=europe-west1 --format="value(apiEndpoint)")
TOKEN=$(gcloud auth print-access-token)
# Executar el pipeline amb macros de data
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 l'estat de les execucions
curl -H "Authorization: Bearer ${TOKEN}" \
"${INSTANCIA}/v3/namespaces/default/apps/erp-costes-a-bigquery/workflows/DataPipelineWorkflow/runs"Els perfils de còmput (compute profiles) defineixen com serà el clúster de Dataproc que executi el pipeline: nombre de workers, tipus de màquina, xarxa, compte de servei. És on es controla el cost d'execució:
Perfil "alpinashop-perfil-computo" Provisioner : Dataproc Region : europe-west1 Master : n2-standard-2, 1 node, disc 100 GB Workers : n2-standard-2, 2 nodes, disc 100 GB Network : alpinashop-vpc / sn-datos-euw1 Internal IP only : true Service Account : [email protected] Image version : 2.2-debian12 Idle TTL : 10 minuts
Aquest perfil aplica exactament el que hem après a 04-03: xarxa privada sense IP pública, compte de servei propi, versió d'imatge fixada i autodestrucció per inactivitat.
- Què passa per sota: Dataproc efímer
Quan prems Run, això és el que passa realment:
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: tradueix el graf visual a un job de Spark
DF->>DP: crea un cluster efimer (2-5 min)
DP->>DP: executa el job de Spark
DP->>BQ: escriu el resultat
DP-->>DF: fi del job
DF->>DP: destrueix el cluster (Idle TTL)
DF-->>U: estat SUCCEEDED + llinatge registrat
Conseqüències molt pràctiques de saber això:
L'arrencada triga de 2 a 5 minuts. No és lentitud de Data Fusion: és el temps de crear el clúster. Per això Data Fusion no serveix per a res que necessiti latència baixa. Un pipeline que mou 200 files triga gairebé el mateix que un que en mou 20 milions, perquè el cost dominant és l'arrencada.
El cost real és la suma de dues coses: les hores d'instància (sempre) més les hores de Dataproc (per execució). Un pipeline diari de 8 minuts amb 3 nodes són uns cèntims de Dataproc, però la instància continua facturant 24 hores al dia.
Pots reutilitzar un clúster. Configurant un perfil que apunti a un clúster existent en comptes de crear-ne un, s'elimina el temps d'arrencada. Té sentit quan s'executen molts pipelines seguits, per exemple en una finestra nocturna: es crea el clúster, es llancen deu pipelines encadenats, es destrueix.
Els errors de Spark apareixen als registres de Dataproc. Quan un pipeline falla per memòria (OutOfMemoryError en un executor) o per biaix de dades, el diagnòstic és exactament el de 04-03: mirar la interfície de Spark i les etapes. Data Fusion no t'amaga aquesta realitat, només l'embolcalla.
- Llinatge de dades a nivell de camp
Aquesta és, per a moltes organitzacions, la raó principal per pagar Data Fusion.
El llinatge respon a dues preguntes que sonen trivials i que gairebé cap empresa no sap contestar:
- "D'on surt exactament aquest camp de l'informe?"
- "Si canvio aquesta columna de l'ERP, què es trenca?"
Data Fusion registra el llinatge automàticament en executar cada pipeline, a dos nivells:
Llinatge de conjunt de dades: quins orígens alimenten quines destinacions, amb quin pipeline i quan.
Llinatge de camp: quina columna concreta de l'origen produeix quina columna de la destinació, i quines operacions ha patit pel camí.
flowchart LR
A["erp_almacen.movimientos_stock<br/>coste_unitario (INT, centims)"]
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
Aquest diagrama no el dibuixa ningú: el genera Data Fusion a partir de l'execució. I respon a la pregunta que a moltes empreses costa dies d'arqueologia: si algú pregunta per què el marge de l'informe de direcció surt estrany, el llinatge mostra que coste_medio_eur ve d'una divisió entre 100 i d'una mitjana no ponderada —que, dit sigui de passada, és una decisió discutible que el llinatge deixa a la vista.
Per què importa tant:
| Escenari | Sense llinatge | Amb llinatge |
|---|---|---|
| Auditoria | "Creiem que ve de l'ERP" | Traçabilitat documentada i datada |
| Canvi a l'origen | Es desplega i es veu què es trenca | Se sap per endavant quins pipelines en depenen |
| Dada sospitosa | Dies d'investigació | Un clic fins a l'origen |
| RGPD: on és una dada personal | Cerca manual per tot arreu | Es rastreja el camp per totes les destinacions |
| Baixa d'un sistema | Por d'apagar-lo | Es veu exactament què en consumeix |
L'última fila és especialment valuosa en una migració: saber què depèn de l'ERP abans de tocar-lo. I la penúltima connecta directament amb el que veurem a 04-07, on Dataplex consolida el llinatge de tota la plataforma, no només el de Data Fusion.
- Replicació amb CDC des del MySQL de l'ERP
Un pipeline per lots que s'executa cada nit deixa les dades amb fins a 24 hores de retard, i a més carrega l'ERP amb una consulta pesada. L'alternativa és la captura de dades de canvi (Change Data Capture, CDC): en comptes de consultar la taula, es llegeix el registre de transaccions de la base de dades i es repliquen els canvis a mesura que passen.
flowchart LR
M["MySQL ERP Sabadell<br/>binlog"]
R["Data Fusion Replication<br/>llegeix el binlog"]
S["Taula de staging<br/>a BigQuery"]
B["Taula desti<br/>alpinashop_analitica.stock_erp"]
M -->|INSERT/UPDATE/DELETE| R
R -->|esdeveniments de canvi| S
S -->|MERGE periodic| B
Preparació al MySQL de l'ERP:
-- Requisits al servidor d'origen (els aplica l'administrador de l'ERP)
-- A 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 és innegociable: amb STATEMENT, el binlog guarda les sentències SQL, no les files resultants, i la replicació no pot reconstruir l'estat. És el primer requisit que cal verificar amb el proveïdor de l'ERP.
A Data Fusion s'utilitza el mòdul Replication (edició Enterprise):
- Origen: MySQL, amb la connexió i l'usuari
cdc_datafusion. - Selecció de taules:
movimientos_stock,proveedores,articulos. - Destinació: BigQuery, dataset
alpinashop_analitica, amb un prefix de staging. - Avaluació de la font: comprova permisos i configuració abans de començar.
- Instantània inicial + streaming continu.
Data Fusion carrega primer una foto completa de les taules i després aplica els canvis en continu, escrivint en taules de staging i executant un MERGE periòdic contra les taules finals.
Consideracions importants sobre CDC, que cal valorar abans de decidir-se:
| Aspecte | Realitat |
|---|---|
| Latència | Segons o pocs minuts, davant de 24 hores del lot |
| Càrrega a l'origen | Molt menor: llegir el binlog no executa consultes |
| Esborrats | Es capturen, cosa que una càrrega incremental per data no fa |
| Cost | Requereix edició Enterprise i execució contínua: és car |
| Requisits | Configuració del servidor d'origen; de vegades el proveïdor no la permet |
| Canvis d'esquema | Es gestionen, però requereixen atenció |
Aquest punt dels esborrats és l'argument tècnic més fort a favor de CDC. Un pipeline incremental que llegeix WHERE fecha_modificacion > X no s'assabenta mai que una fila s'ha eliminat, i el magatzem analític acumula registres fantasma indefinidament. CDC sí que veu el DELETE.
L'alternativa a considerar és Datastream, un servei de Google dedicat exclusivament a CDC, més simple i més barat que activar Enterprise a Data Fusion només per a això. És a la taula de l'apartat següent.
- Quan Data Fusion, quan Dataflow, quan Datastream, quan
bq load
bq loadLa taula honesta, que és el que cal endur-se d'aquesta lliçó:
| Criteri | bq load |
Datastream | Data Fusion | Dataflow |
|---|---|---|---|---|
| Cas típic | Fitxer llest → taula | Rèplica de BD amb CDC | Integrar orígens variats amb neteja | Transformació a mida, streaming |
| Codi | Una comanda | Cap | Cap (visual) | Python o Java |
| Perfil necessari | Qualsevol | Qualsevol + DBA de l'origen | Analista | Enginyer de dades |
| Cost fix | Zero | Per GB processat | Per hora d'instància | Zero en lot |
| Cost variable | Gratis | Baix | Dataproc per execució | Treballadors per hora |
| Latència mínima | Minuts | Segons | 2-5 min d'arrencada | Segons (streaming) |
| Transformacions | Cap | Cap (replica tal qual) | Moltes, visuals | Il·limitades |
| Streaming real | No | Sí (replicació) | Limitat | Sí, el seu punt fort |
| Connectors externs | No | Bases de dades | Moltíssims | Els que programis |
| Llinatge | No | No | Sí, per camp | Via Dataplex |
| Versionat a Git | Trivial | N/A | Incòmode (JSON exportat) | Trivial |
I traduït a regles de decisió:
Utilitza bq load quan la dada ja és a Cloud Storage amb la forma correcta. És gratis i no hi ha res a mantenir. Comença sempre preguntant-te si això n'hi ha prou. La meitat dels pipelines que existeixen al món sobren.
Utilitza Datastream quan l'objectiu sigui replicar una base de dades completa a BigQuery amb latència baixa i captura d'esborrats, sense transformar res. És més simple i barat que Data Fusion Enterprise per a aquesta tasca concreta, i admet MySQL, PostgreSQL, Oracle i SQL Server.
Utilitza Data Fusion quan hi hagi diversos orígens heterogenis, calguin transformacions i neteja no trivials, l'equip no tingui perfil de programació, i el llinatge sigui un requisit real (auditoria, compliment normatiu). I quan el volum de feina justifiqui el cost de la instància.
Utilitza Dataflow quan hi hagi streaming amb semàntica de temps de l'esdeveniment, transformacions complexes o dependents d'estat, o quan l'equip prefereixi codi versionable amb proves automàtiques.
I hi ha una combinació molt raonable que es veu molt a la pràctica: Datastream per replicar el cru + BigQuery SQL per transformar (ELT pur). Sense Data Fusion, sense Dataflow, sense instàncies que pagar. Per a molts casos, inclosa bona part del d'AlpinaShop, és la resposta més eficient.
- La decisió raonada d'AlpinaShop
Amb tot l'anterior sobre la taula, la Lucía i la Marta decideixen així:
L'ERP de Sabadell → Datastream, no Data Fusion. El requisit és replicar tres taules amb latència baixa i capturar esborrats. No hi ha transformació: els costos es calculen després amb SQL a BigQuery, on la Lucía es mou perfectament. Datastream costa una fracció d'una instància Enterprise i no hi ha instància que apagar. Les transformacions de l'apartat 8 (majúscules, cèntims a euros, etiquetes de tipus de moviment) es converteixen en una vista de BigQuery, que a més queda versionada a Git.
El CSV mensual del transportista → Data Fusion, amb instància efímera. Aquí sí que guanya: el fitxer és brut de veritat i les onze directives de Wrangler s'escriuen en vint minuts veient les dades, davant d'un parell de dies de desenvolupament i proves en Beam. És mensual, així que la instància s'encén, s'executa i s'apaga. Amb edició Basic i menys de 120 hores al mes, el cost és zero. La instància es gestiona així:
# Apagar la instancia quan no s'utilitza (Enterprise i Basic ho permeten)
gcloud data-fusion instances update alpinashop-fusion \
--location=europe-west1 --enable-instance-stopped
# I tornar a encendre-la el dia de la carrega mensual
gcloud data-fusion instances update alpinashop-fusion \
--location=europe-west1 --no-enable-instance-stoppedL'API del proveïdor de motxilles → una Cloud Function programada. Són 400 productes en JSON, un cop a la setmana. Aixecar un clúster de Dataproc de tres nodes durant cinc minuts per llegir 400 registres és desproporcionat. Quaranta línies de Python a Cloud Functions (06-03) disparades per Cloud Scheduler resolen el cas per cèntims.
La regla general que AlpinaShop adopta, i que serveix com a criteri reutilitzable:
- La dada ja és a Cloud Storage amb forma de taula? →
bq load. - Cal replicar una base de dades sense transformar? → Datastream + vistes de BigQuery.
- És un volum petit d'una API? → Cloud Function programada.
- És un fitxer brut, recurrent, que netejarà un analista? → Data Fusion amb instància efímera.
- Hi ha streaming, finestres o lògica complexa? → Dataflow.
- És algorítmic amb biblioteques de Spark? → Dataproc Serverless.
Aquesta llista, amb la disciplina de recórrer-la en ordre i quedar-se a la primera que serveixi, és el que impedeix que una plataforma de dades d'una pime acabi costant el que la d'una multinacional.
Errors Habituals i Consells
Deixar la instància encesa. És l'error car d'aquest servei i mereix repetir-se: una instància Basic oblidada costa més de mil euros al mes sense executar ni un sol pipeline. Apaga-la o esborra-la.
Triar Enterprise sense necessitar-ho. Només cal per a CDC, alta concurrència i alta disponibilitat. Basic cobreix la majoria dels casos d'una pime a menys de la meitat de preu.
Utilitzar Data Fusion per a volums minúsculs. Un clúster de Dataproc per processar 400 files és desproporcionat. L'arrencada triga més que la feina.
No versionar els pipelines. Encara que siguin visuals, s'han d'exportar a JSON i guardar a Git. Sense això, una instància esborrada s'endú mesos de feina i no hi ha historial de canvis.
# Exportar un pipeline per versionar-lo
curl -H "Authorization: Bearer $(gcloud auth print-access-token)" \
"${INSTANCIA}/v3/namespaces/default/apps/erp-costes-a-bigquery" \
> pipelines/erp-costes-a-bigquery.jsonParal·lelitzar l'extracció sense mesurar l'impacte a l'origen. Number of Splits = 8 sobre el MySQL de l'ERP pot deixar sense servei la bàscula del magatzem. Comença per 1 i puja mesurant.
Escriure credencials a la configuració del pipeline. Utilitza macros resoltes des de Secret Manager. Un pipeline exportat a JSON amb la contrasenya dins acaba a Git, i això és un incident de seguretat.
Oblidar la paginació al plugin HTTP. Es porten els primers 100 registres i ningú no ho nota fins que falten productes en un informe.
Saltar-se el Preview. És gratis, triga segons i mostra exactament què surt de cada bloc. Desplegar sense previsualitzar és llançar un clúster per descobrir un error de tipus.
No utilitzar Upsert a la destinació. Amb Insert, cada reexecució duplica files. Amb Upsert sobre la clau de negoci, el pipeline és idempotent i es pot rellançar sense por.
Consell: utilitza macros des del principi. Dates, rutes, noms de taula i credencials com a ${variable}. Converteix un pipeline rígid en un de reutilitzable i orquestrable.
Consell: posa sempre un Error Collector. Els registres rebutjats han d'anar a algun lloc amb el seu motiu. És la mateixa disciplina que la quarantena de Beam a 04-02.
Consell: revisa el llinatge després de cada desplegament. És la millor manera de verificar que el pipeline fa el que creus que fa, i no costa res.
Exercicis
Exercici 1: directives de Wrangler per al fitxer del proveïdor
El proveïdor de motxilles envia aquest fitxer de text separat per tabuladors:
REF DESCRIPCION PVP_RECOMENDADO COSTE STOCK ALTA ACTIVO mb-4001 Motxilla Trekking 40L Blava 89,90 EUR 52,30 EUR 120 01-03-2024 S MB-4002 motxilla trekking 30l vermella 74,50 EUR 43,10 EUR 0 15-06-2025 S MB-4003 Motxilla Alpina 55L 129,00 EUR 78,00 EUR - N/D N
Escriu les directives de Wrangler que produeixin un conjunt net amb: sku en majúscules sense espais; nombre amb format de títol; pvp_eur i coste_eur com a números decimals sense la paraula EUR; stock com a enter, amb - convertit en 0; fecha_alta com a data yyyy-MM-dd, tolerant N/D com a nul; activo com a booleà; i una columna calculada margen_pct. Els registres sense REF s'han de desviar a error.
Exercici 2: dissenyar el pipeline d'enviaments
Dissenya el pipeline complet que llegeix el CSV mensual del transportista des de gs://alpinashop-datalake/transportista/2026/03/envios.csv, el neteja, el creua amb la taula pedidos d'alpinashop_analitica per afegir el país i el canal, calcula el retard en dies entre fecha_pedido i fecha_entrega, i escriu a alpinashop_analitica.envios. Indica: cada bloc amb el seu tipus i configuració essencial, com tractar els enviaments el pedido_id dels quals no existeixi a pedidos, quina operació utilitzar a la destinació i per què, i quines macros definiries.
Exercici 3: la decisió d'eina, amb justificació econòmica
AlpinaShop absorbeix un competidor i hereta quatre integracions noves. Per a cadascuna, tria entre bq load, Datastream, Data Fusion, Dataflow, Dataproc Serverless o Cloud Function, i justifica-ho incloent-hi una estimació de l'ordre de magnitud del cost mensual:
- Un PostgreSQL 14 de 80 GB amb l'històric de clients del competidor, que s'ha de replicar a BigQuery amb menys de 5 minuts de retard i capturant esborrats.
- Un fitxer Excel setmanal de 3.000 files que envia l'equip de compres, amb columnes que canvien de nom cada dos per tres i dades escrites a mà.
- Un flux de 2.000 esdeveniments per segon d'una aplicació mòbil que cal agregar per finestres de 5 minuts abans de desar-lo.
- Un bolcat nocturn en Parquet de 12 GB que el competidor ja deixa en un bucket, amb l'esquema exacte de la taula de destinació.
Solucions
Solució 1
-- 1) Parsejar el fitxer separat per tabuladors, amb capcalera parse-as-csv :body '\t' true drop :body -- 2) SKU: treure espais i normalitzar a majuscules trim :REF uppercase :REF rename :REF sku -- 3) Desviar a error els registres sense referencia send-to-error sku == null || sku.isEmpty() -- 4) Nom: netejar espais i aplicar format de titol trim :DESCRIPCION titlecase :DESCRIPCION rename :DESCRIPCION nombre -- 5) Imports: treure ' EUR', canviar la coma decimal i 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 zero, no nul find-and-replace :STOCK s/^-$/0/g set-type :STOCK int rename :STOCK stock -- 7) Data: 'N/D' a nul, i format dd-MM-yyyy a data 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) Actiu: 'S'/'N' a boolea set-column :ACTIVO (ACTIVO == 'S') rename :ACTIVO activo -- 9) Columna calculada: marge percentual, protegit de la divisio per zero set-column :margen_pct (pvp_eur != null && pvp_eur > 0) ? ((pvp_eur - coste_eur) / pvp_eur * 100) : null
Resultat esperat sobre les tres files d'exemple:
| sku | nombre | pvp_eur | coste_eur | stock | fecha_alta | activo | margen_pct |
|---|---|---|---|---|---|---|---|
| MB-4001 | Motxilla Trekking 40L Blava | 89.90 | 52.30 | 120 | 2024-03-01 | true | 41.8 |
| MB-4002 | Motxilla Trekking 30L Vermella | 74.50 | 43.10 | 0 | 2025-06-15 | true | 42.1 |
| MB-4003 | Motxilla Alpina 55L | 129.00 | 78.00 | 0 | null | false | 39.5 |
Els tres punts que s'avaluen: distingir - (que significa zero unitats) de N/D (que significa dada desconeguda, és a dir, nul) —confondre'ls falsejaria qualsevol informe d'estoc—; protegir la divisió de la columna calculada; i desviar a error en comptes de filtrar en silenci, perquè un fitxer amb referències buides deixi rastre.
Solució 2
Bloc 1 — Source: GCS
- Path:
gs://alpinashop-datalake/transportista/${anyo}/${mes}/envios.csv - Format:
text(una fila per línia; el parseig es fa a Wrangler perquè el separador és;) - Service Account:
sa-fusion-pipelines@...
Bloc 2 — Transform: Wrangler
Les directives de l'apartat 7, incloent-hi el drop de nom i cognoms per minimització de dades, i send-to-error per a les files sense pedido_id.
Bloc 3 — Source: BigQuery
- Dataset/Table:
alpinashop_analitica.pedidos - Import Query (millor que llegir la taula sencera):
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 filtre de partició és obligatori: la taula té require_partition_filter=TRUE des de 04-01, així que sense ell el pipeline fallaria. I encara que no el tingués, llegir la taula sencera cada mes seria llençar diners.
Bloc 4 — Analytics: Joiner
- Inputs: sortida de Wrangler (esquerra) i sortida de BigQuery (dreta)
- Join type:
Left Outersobrepedido_id - Fields: tots els del CSV net, més
fecha_pedido,paisicanaldel costat dret
El tipus d'unió és la decisió clau de l'exercici. Amb Inner, els enviaments el pedido_id dels quals no existeixi a pedidos desapareixerien sense deixar rastre, i ningú no sabria que falten enviaments. Amb Left Outer es conserven tots, amb pais i canal nuls, i aquests nuls són el senyal que hi ha un problema d'integritat a investigar: comandes de mesos fora del rang de dates, o identificadors mal escrits al fitxer del transportista.
Bloc 5 — Transform: Wrangler (segon)
-- Retard en dies entre comanda i lliurament set-column :retraso_dias (fecha_entrega != null && fecha_pedido != null) ? (dateDiff(fecha_entrega, fecha_pedido)) : null -- Marca d'enviament orfe, per poder comptar-los en un informe de qualitat set-column :sin_pedido_asociado (pais == null)
Bloc 6 — Sink: BigQuery
- Table:
alpinashop_analitica.envios - Operation:
Upsertamb clauenvio_id - Partition field:
fecha_entrega; Cluster field:pais
Per què Upsert. El fitxer del transportista se sol reenviar corregit quan hi ha errors, i de vegades inclou enviaments del mes anterior que es van lliurar tard. Amb Insert, cada reenviament duplicaria files i l'informe de terminis mentiria. Amb Upsert sobre envio_id, reexecutar el pipeline les vegades que calgui produeix sempre el mateix resultat: és idempotent, exactament el mateix principi que exigíem als consumidors de Pub/Sub a 04-04.
Bloc 7 — Error Collector → Sink GCS
- Path:
gs://alpinashop-datalake/cuarentena/envios/${anyo}/${mes}/ - Format:
json
Macros a definir: ${anyo}, ${mes}, ${fecha_desde}, ${fecha_hasta}. Amb elles, el mateix pipeline serveix per a qualsevol mes i es pot reprocessar un històric complet canviant només els arguments. Sense elles, caldria duplicar el pipeline o editar-lo cada mes.
Solució 3
1. PostgreSQL de 80 GB replicat a BigQuery, menys de 5 minuts de retard, amb esborrats → Datastream.
És literalment la seva definició: CDC gestionat sobre PostgreSQL, sense codi, amb captura de DELETE —que una càrrega incremental per data no detectaria mai—. Data Fusion Enterprise també podria, però exigiria l'edició cara (~3.000 $/mes d'instància) per fer exactament el mateix. Datastream es factura per GB processat: la càrrega inicial de 80 GB més els canvis diaris situen el cost en l'ordre d'unes desenes d'euros al mes. La transformació posterior es fa amb vistes a BigQuery, gratis en esforç i versionables a Git.
2. Excel setmanal de 3.000 files amb columnes canviants i dades a mà → Data Fusion amb instància efímera. És el cas canònic de Wrangler: dades brutes escrites per persones, esquema inestable, i un analista —no un programador— que necessita corregir-ho veient les dades. Programar-ho en Beam significaria redesplegar codi cada vegada que compres reanomeni una columna. Amb instància Basic encesa una hora a la setmana, són 4 hores al mes, molt dins de les 120 gratuïtes: cost efectiu zero, més uns cèntims de Dataproc per execució. La clau de la resposta és la paraula efímera: si la instància es deixa encesa, la mateixa solució costa més de mil euros al mes.
3. 2.000 esdeveniments per segon agregats en finestres de 5 minuts → Dataflow.
Streaming amb finestres temporals: el territori exclusiu de Beam. Data Fusion no fa això bé, bq load no aplica i una Cloud Function no pot mantenir estat de finestra. Un pipeline de streaming amb 2-3 treballadors permanents està en l'ordre de 100-150 € al mes, cost que cal assumir conscientment perquè és l'única eina que resol el requisit. Si el negoci tolerés 15 minuts de retard en comptes de 5, l'alternativa —subscripció de Pub/Sub a Cloud Storage més micro-lots— costaria una fracció; val la pena preguntar-ho abans d'encendre el pipeline.
4. Parquet nocturn de 12 GB amb l'esquema exacte de la taula de destinació → bq load.
Sense transformació i amb l'esquema ja correcte, no hi ha res a processar. La càrrega per lots a BigQuery és gratuïta i Parquet porta l'esquema incorporat, així que ni tan sols cal declarar-lo. Una comanda a l'orquestrador de 04-06, disparada per la notificació del bucket que vam muntar a 04-04. Cost: 0 € de procés, només l'emmagatzematge dels 12 GB. Qualsevol altra opció d'aquesta llista seria pagar per fer una còpia, i és exactament el reflex que cal corregir: la primera pregunta davant de qualsevol integració és "n'hi ha prou amb bq load?".
Conclusió
Has vist el costat menys vistós i més real d'una plataforma de dades: les dades que no neixen al núvol. El MySQL de l'ERP a Sabadell, el CSV del transportista amb les seves dates europees i els seus nuls escrits com a guionet, i l'API del proveïdor de motxilles. Cap no publica a Pub/Sub, cap no escriu Parquet, i la Lucía necessita els tres per calcular el marge real per producte i explicar per què cauen les ressenyes en certes zones.
Has entès què és un ETL/ELT visual i, sobretot, a qui serveix: a perfils que dominen el negoci i el SQL però no escriuran Beam, i a organitzacions que necessiten traçabilitat documentada. I també a qui no serveix, que és igual d'important. Saps que Data Fusion és CDAP gestionat, amb el seu Studio, el seu Wrangler, el seu Hub i el seu registre de llinatge, i que per sota no executa res: tradueix el graf visual a Spark i aixeca un Dataproc efímer, cosa que explica els seus 2-5 minuts d'arrencada i per què no serveix per a latències baixes.
Has posat el cost al davant en comptes d'amagar-lo: edicions Developer, Basic i Enterprise, facturades per hora d'instància existent, s'utilitzi o no, amb una instància Basic permanent costant més que tota la resta de la plataforma d'AlpinaShop junta. Aquesta dada no desqualifica el producte —per a una empresa amb vint orígens és barat davant dels sous que estalvia— però sí que obliga a l'estratègia d'instància efímera en una pime.
Has netejat el CSV del transportista amb Wrangler i les seves directives: parseig per punt i coma, trim, nuls escrits com a guionet, dates dd/MM/yyyy convertides a tipus data, decimals amb coma, divisió de la columna de destinatari, format de títol, columnes derivades, i el drop final del nom i els cognoms per minimització de dades, amb l'advertiment exprés de RGPD i de revisió per compliance. Has valorat l'avantatge real de l'eina, que no és la potència sinó el cicle de retroalimentació: veus l'efecte de cada directiva sobre dades reals a l'instant.
Has dissenyat el pipeline de l'ERP amb les seves macros de data, la seva lectura particionada amb la prudència de no tombar la bàscula del magatzem, la seva neteja, la seva agregació i la seva destinació en Upsert perquè sigui idempotent i reexecutable. Coneixes el Hub i els seus connectors, el plugin HTTP amb la seva paginació fàcil d'oblidar, el flux Preview → Deploy → Run → Schedule, els perfils de còmput que controlen el clúster subjacent, i l'API REST amb la qual l'orquestrador l'invocarà. Has vist el llinatge a nivell de camp, que respon a "d'on surt aquesta columna?" i "què es trenca si canvio això?" sense dies d'arqueologia, i la replicació CDC des de MySQL amb el seu binlog_format = ROW innegociable i el seu gran argument: és l'única manera d'assabentar-se dels esborrats.
I has tancat amb la taula honesta i amb la decisió raonada d'AlpinaShop, que no és "fem servir Data Fusion per a tot" sinó una llista ordenada: primer bq load si n'hi ha prou, després Datastream si es tracta de replicar, després una Cloud Function si és poca cosa, després Data Fusion amb instància efímera si el fitxer és brut i el neteja un analista, i Dataflow o Dataproc si hi ha streaming o algorismes. Recórrer aquesta llista en ordre i quedar-se a la primera opció que serveixi és el que separa una plataforma de dades proporcionada d'una caríssima.
I ara apareix un problema nou, que és conseqüència directa de l'èxit de les quatre lliçons anteriors. AlpinaShop té una exportació nocturna de Cloud SQL, una càrrega a BigQuery, un pipeline de Dataflow, una feina de Spark a Dataproc Serverless, un pipeline mensual de Data Fusion, diverses consultes d'agregació i una vista materialitzada per refrescar. Són set o vuit processos que depenen els uns dels altres: no té sentit llançar l'agregació abans que hagi acabat la càrrega, ni refrescar la vista de negoci amb dades a mitges.
Ara mateix, això ho resol un cron en una màquina virtual que llança scripts a hores fixes, calculades a ull amb marge de sobres. Funciona fins a la primera nit en què l'exportació triga vint minuts més del normal: llavors la càrrega llegeix un fitxer incomplet, l'agregació calcula sobre dades parcials, l'informe de direcció es lleva amb xifres falses, i ningú no se n'assabenta fins que algú les mira a mig matí.
A 04-06, Cloud Composer i Workflows, resoldrem això. Veuràs per què un cron no és un orquestrador i què significa realment coordinar dependències, reintents, alertes i reprocessos. Coneixeràs Cloud Composer —Apache Airflow gestionat— amb els seus DAG, tasques, operadors i sensors, i escriuràs el DAG complet del procés nocturn d'AlpinaShop comentat línia a línia, amb l'advertiment clar que Composer és car per a una pime. I coneixeràs Workflows, l'orquestració serverless en YAML sense cost fix, amb Cloud Scheduler per al dispar, que és per on AlpinaShop començarà.
Curs de Google Cloud Platform (GCP)
Mòdul 1: Introducció a Google Cloud Platform
- Què és Google Cloud Platform?
- Configuració del teu compte de GCP
- Descripció general de la consola de GCP
- Projectes, jerarquia de recursos i facturació
- Regions, zones i model de responsabilitat compartida
- Cloud Shell i la CLI de gcloud
Mòdul 2: Serveis principals de GCP
- Compute Engine: màquines virtuals a Google Cloud
- Cloud Storage: emmagatzematge d'objectes
- Cloud SQL: bases de dades relacionals gestionades
- App Engine: plataforma com a servei
- Google Kubernetes Engine (GKE)
- Bases de dades NoSQL: Firestore, Bigtable i Spanner
- Com triar el servei de còmput adequat
Mòdul 3: Xarxes i seguretat
- Xarxes VPC
- Balanceig de càrrega al núvol
- Cloud CDN
- Gestió d'identitat i accés (IAM)
- Cloud Armor
- Secrets i xifratge: Secret Manager i Cloud KMS
- Cloud DNS, certificats TLS i publicació segura de serveis
Mòdul 4: Dades i anàlisi
- BigQuery: el magatzem de dades analític
- Cloud Dataflow: processament de dades per lots i en temps real
- Cloud Dataproc: Spark i Hadoop gestionats
- Cloud Pub/Sub: missatgeria asíncrona
- Cloud Data Fusion: integració de dades sense codi
- Orquestració de pipelines amb Cloud Composer i Workflows
- Govern de les dades i taulers amb Dataplex i Looker Studio
Mòdul 5: Aprenentatge automàtic i IA
- Vertex AI: la plataforma d'aprenentatge automàtic de GCP
- AutoML: models a mida sense escriure codi
- TensorFlow a GCP: entrenament i servei de models
- API de llenguatge natural
- API de visió
- IA generativa a Vertex AI: models Gemini i incrustacions
- MLOps: del model al producte amb Vertex AI Pipelines
Mòdul 6: DevOps i monitoratge
- Cloud Build: integració contínua a GCP
- Cloud Source Repositories i gestió del codi font
- Cloud Functions: funcions sense servidor
- Cloud Monitoring (abans Stackdriver): mètriques, taulers i alertes
- Cloud Deployment Manager i infraestructura com a codi nativa
- Cloud Logging i Cloud Trace: registres, traces i diagnòstic
- Terraform a GCP: infraestructura com a codi a la pràctica
Mòdul 7: Temes avançats de GCP
- Híbrid i multinúvol amb Anthos
- Computació sense servidor amb Cloud Run
- Xarxes avançades: VPC compartida, aparellament i connectivitat híbrida
- Bones pràctiques de seguretat
- Gestió i optimització de costos
- Fiabilitat: SLO, alta disponibilitat i recuperació de desastres
- Govern a escala: organització, polítiques i auditoria
