Everything we have built so far shares one comfortable assumption: that the data is already inside Google Cloud. The history came from Cloud SQL, the events arrive over Pub/Sub, the exports land in alpinashop-catalogo. But AlpinaShop, like any real company, has data that lives elsewhere and arrives by other means.
There are three concrete cases, and none of them is exotic:
- The warehouse ERP is a MySQL 8 running on a physical server in the Sabadell offices. It holds the real stock, the goods received, the suppliers and the purchase costs. Nobody is going to migrate it this year: it works, it is wired up to the weighing scale and to the label printer, and the vendor charges for every change.
- The carrier emails a CSV every month with the shipments, the incidents and the real delivery times. Semicolon-separated, dates in
dd/mm/yyyyformat, amounts with a decimal comma, city names in capitals with no accents, and aDESTINATARIOcolumn that mixes first name and surname. - The backpack supplier offers a REST API with its up-to-date catalogue: references, cost prices, availability and technical specifications.
Lucía needs all three. She needs to cross the ERP's purchase costs with the sales in BigQuery to work out the real margin per product. She needs the carrier's delivery times to explain why customer reviews are dropping in certain areas. And she needs the supplier's catalogue to spot discontinued products that are still published on the website.
Lucía knows SQL. She knows a great deal of SQL. But she has never written an Apache Beam pipeline in her life and is not about to start, and Dani has his own to-do list. If every new file requires a development project, that data will never arrive.
Cloud Data Fusion exists for exactly that gap: a visual data integration tool where pipelines are built by dragging blocks around and data is cleansed while you look at it, without writing code. In this lesson you will use it to build the ERP pipeline, you will cleanse the carrier's CSV with Wrangler, you will understand the field-level lineage that is its greatest virtue, you will set up replication with change data capture from MySQL, and — this matters just as much — you will learn when not to use it, because it carries an hourly cost that shapes the entire decision.
Contents
- The problem of data that is not born in Google Cloud
- What a visual ETL/ELT is and who it serves
- What Data Fusion is: managed CDAP
- Editions, hourly cost and its consequences
- Creating the instance and understanding what has been created
- The Studio: sources, transformations and sinks
- Wrangler: cleansing the carrier's CSV while looking at the data
- The real pipeline: the ERP's MySQL into
alpinashop_analitica - Connectors and plugins from the Hub
- Deploying, running and scheduling
- What happens underneath: ephemeral Dataproc
- Field-level data lineage
- Replication with CDC from the ERP's MySQL
- When Data Fusion, when Dataflow, when Datastream, when
bq load - AlpinaShop's reasoned decision
- The problem of data that is not born in Google Cloud
This is called data integration, and it is the least glamorous and most time-consuming work in any analytics project. Industry surveys repeat it without variation: between 60% and 80% of the effort in a data project goes into getting the data to arrive, not into analysing it.
The concrete problems, with AlpinaShop's examples:
| Problem | Case at AlpinaShop |
|---|---|
| Connectivity | The MySQL sits in Sabadell, behind a firewall, with no public IP |
| Formats | CSV with ;, dd/mm/yyyy dates, decimal commas, no reliable header |
| Quality | Cities in capitals without accents, nulls written as -, stray whitespace |
| Changing schemas | The carrier added a column in January without warning |
| Uneven frequency | The ERP changes continuously; the CSV arrives monthly; the API is queried on demand |
| Volume | The ERP's history is 12 million stock movements |
| Traceability | Nobody knows where the coste_unitario field in a 2025 report came from |
Every one of these problems can be solved by writing code. The trouble is that solving them by writing code for each source, and maintaining that code when the carrier changes the format, is continuous work that a small business of 40 people cannot afford to hand to Dani.
- What a visual ETL/ELT is and who it serves
ETL stands for extract, transform, load: the data is pulled out of the source, transformed outside, and loaded into the destination already clean. ELT swaps the last two around: the data is loaded raw and transformed inside the destination, taking advantage of its power.
| Approach | Where the transformation happens | Advantage | Drawback |
|---|---|---|---|
| ETL | In an intermediate engine | The destination only ever receives clean data | That engine has to be sized and paid for |
| ELT | In the destination (BigQuery) | Uses its engine; you keep the raw data | The destination stores dirty data; query cost |
With modern warehouses such as BigQuery, the clear trend is ELT: load raw and transform with SQL. But there is one part that remains unavoidably ETL: getting the data out of the source and up to the door. That is the extraction, and it is exactly where Data Fusion earns its keep.
A visual ETL is a tool where the pipeline is built on a canvas with blocks instead of with code. Who does it serve?
It serves very well:
- Analysts like Lucía, who know the business and the data but do not program distributed pipelines.
- Small teams with no dedicated data engineers.
- Standard integrations: read a table, clean a few fields, write into another table.
- Organizations that need documented traceability of where each piece of data comes from (audits, compliance).
It serves badly:
- Complex business logic with nested conditions and state.
- Teams that already have engineers and mature version control: a visual pipeline versions far worse than a
.pyfile. - Streaming with windows and event time: that is what Beam is for.
- Tight budgets with sporadic use, for the reason set out in section 4.
And an honest warning worth making from the start: code-free does not mean knowledge-free. To build a pipeline in Data Fusion you have to understand schemas, types, joins, keys and partitions exactly as you would if you were programming it. What you save is the syntax and the infrastructure, not the thinking.
- What Data Fusion is: managed CDAP
Cloud Data Fusion is the managed version of CDAP (Cask Data Application Platform), an open source data integration platform that Google acquired and offers as a service.
Its components:
| Component | What it does |
|---|---|
| Studio | The visual canvas where the pipeline is drawn |
| Wrangler | Interactive data explorer and cleanser |
| Hub | Catalogue of plugins, connectors and sample pipelines |
| Metadata and lineage | Automatic record of which field comes from where |
| Replication | CDC module for copying databases continuously |
| Execution engine | Generates and runs the job on Dataproc |
That last row is the most important one for understanding the product: Data Fusion does not execute anything itself. It translates the visual pipeline into a Spark or MapReduce job and runs it on a Dataproc cluster that it creates on the fly. That explains its behaviour, its start-up time and a good part of its cost, and we will develop it in section 11.
Being managed CDAP has another relevant consequence: the pipelines are portable. A pipeline exported from Data Fusion can be imported into a self-managed CDAP anywhere. It is not a closed proprietary format.
- Editions, hourly cost and its consequences
Here is the characteristic that conditions every decision about this product, and it deserves to go up front rather than be hidden away at the end.
| Edition | What for | Approximate cost per instance hour | Notes |
|---|---|---|---|
| Developer | Testing and development | ~$0.35 | No high availability, limited capacity |
| Basic | Simple production | ~$1.80 | Includes 120 free hours a month per account |
| Enterprise | Demanding production | ~$4.20 | High availability, more concurrency, full lineage, CDC |
Check the current prices in the official documentation; what matters is the model, not the exact figure.
And the model is this: you pay per hour that the instance exists, whether you use it or not. Not per pipeline run, nor per byte processed. The instance is an environment that is switched on.
Let us do the sums for AlpinaShop, because this is the conversation you have to have with management:
| Scenario | Hours/month | Approximate cost |
|---|---|---|
| Permanent Basic instance | 730 | ~$1,310 (less 120 free hours: ~$1,100) |
| Permanent Enterprise instance | 730 | ~$3,070 |
| Permanent Developer instance | 730 | ~$255 |
| Basic switched on 4 h a day | 120 | $0 (within the 120 free hours) |
And on top of that you have to add the cost of the Dataproc cluster that is spun up to run each pipeline, which in practice is usually smaller but is not zero.
The consequence is direct and has to be said without decoration: a permanent Data Fusion instance costs more per month than the whole of the rest of AlpinaShop's data platform put together. BigQuery costs a few euros, Pub/Sub cents, batch Dataflow cents. A permanent Data Fusion Basic would cost over a thousand euros.
That does not invalidate the product: for a company with twenty data sources and an integration team, a thousand euros a month is a bargain against the salaries it saves. But for a small business with three sources the equation does not add up, and that has to be said.
The strategy that does work in a small business is to treat the instance the way we treated the Dataproc clusters in 04-03: as ephemeral. Switch it on to develop, switch it off when you are done; and the pipelines that are already deployed run all the same, because the work is done by Dataproc. We will come back to this in section 15.
- Creating the instance and understanding what has been created
gcloud config set project alpinashop-datos
gcloud services enable datafusion.googleapis.com
# Development instance: the cheapest one for learning and building
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=analiticaCreation takes between 15 and 25 minutes. That is not a fault: a complete CDAP environment is being provisioned in a Google-managed project.
That "managed project" detail matters for permissions. Data Fusion creates the instance in a Google-owned project (the tenant project) and acts on yours through a service agent, which has to be granted permissions explicitly:
PROJ_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_FUSION="service-${PROJ_NUM}@gcp-sa-datafusion.iam.gserviceaccount.com"
# The agent needs to be able to act on the project
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:${SA_FUSION}" \
--role="roles/datafusion.serviceAgent"
# Service account used by the Dataproc clusters that run the pipelines
gcloud iam service-accounts create sa-fusion-pipelines \
--display-name="Data Fusion pipeline execution"
SA_PIPE="[email protected]"
for ROLE 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="$ROLE"
done
# The Data Fusion agent must be able to use that account
gcloud iam service-accounts add-iam-policy-binding "$SA_PIPE" \
--member="serviceAccount:${SA_FUSION}" \
--role="roles/iam.serviceAccountUser"Access to the interface:
gcloud data-fusion instances describe alpinashop-fusion \
--location=europe-west1 --format="value(apiEndpoint, serviceEndpoint, state)"And connectivity with Sabadell, which is the real requirement for reading from the ERP. There are three options, in order of preference:
| Option | How it works | Assessment |
|---|---|---|
| Cloud VPN or Interconnect | Tunnel between alpinashop-vpc and the office network |
The right one. The MySQL is never exposed to the internet |
| Public IP with an allow list and TLS | The port is opened only to specific ranges | Grudgingly acceptable; unnecessary attack surface |
| Export to a file and upload it | A script on the ERP dumps CSV into Cloud Storage | Simple, but rules out CDC and fresh data |
For AlpinaShop, Marta sets up a Cloud VPN tunnel with routing towards the ERP subnet, consistent with everything covered in 03-01: the MySQL still has no public IP and the traffic travels encrypted. The detailed configuration of hybrid connectivity is 07-03's territory.
- The Studio: sources, transformations and sinks
The Studio is a canvas. On the left there is a palette of blocks grouped by category; you drag them onto the canvas and connect them with arrows that represent the flow of the data.
The blocks are classified like this:
| Category | What it does | Examples |
|---|---|---|
| Source | Reads data | Database (JDBC), BigQuery, GCS, Salesforce, HTTP, Kafka |
| Transform | Modifies records | Wrangler, JavaScript, Python, Projection, Encoder |
| Analytics | Aggregates and joins | Group By, Joiner, Deduplicate, Distinct, Row Denormalizer |
| Conditions and Actions | Flow control and actions | Condition, Email, BigQuery Execute, Database Execute |
| Sink | Writes data | BigQuery, GCS, Database, Spanner, Pub/Sub |
| Error Handlers | Collect rejected records | Error Collector |
Every block has a configuration panel: connection string, table, output schema, specific options. And every block declares an output schema — the list of fields with their types — which is propagated to the next one. That propagation is what makes the lineage of section 12 possible.
The pipeline we are going to build, drawn out:
flowchart LR
S1["Source: Database<br/>MySQL ERP Sabadell<br/>table movimientos_stock"]
S2["Source: GCS<br/>carrier CSV"]
W1["Transform: Wrangler<br/>cleansing and types"]
W2["Transform: Wrangler<br/>dates, decimals, names"]
J["Analytics: Joiner<br/>by sku"]
G["Analytics: Group By<br/>average cost per sku"]
K1["Sink: BigQuery<br/>alpinashop_analitica.stock_erp"]
K2["Sink: BigQuery<br/>alpinashop_analitica.envios"]
E["Error Collector<br/>-> GCS quarantine"]
S1 --> W1 --> G --> K1
S2 --> W2 --> J
W1 --> J
J --> K2
W2 -.rejects.-> E
Notice two things in the diagram, because they reproduce the good practices of the previous lessons:
- The rejects branch exists here too. The
Error Collectorblock picks up the records that aWranglercould not process and sends them to a quarantine destination, exactly like Beam's tagged outputs in 04-02. - A single source feeds two branches. The canvas is a directed graph, not a line, just like Beam's graph.
- Wrangler: cleansing the carrier's CSV while looking at the data
Wrangler is the piece that on its own justifies learning Data Fusion. It is an interactive environment where you load a sample of the data, see it as a table, and apply transformations — called directives — while watching the effect immediately.
We start from the carrier's real CSV, exactly as it arrives:
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;DELAY 24H;PED-2026-0043 ENV0098213;-;FERNANDEZ SANZ, ANA;MADRID;28004;5,000;12,00;WRONG ADDRESS;PED-2026-0044
The problems are obvious at a glance: ; separator, European dates, decimal commas, nulls as -, stray whitespace, inconsistent capitalisation, first name and surname bundled together, and a delivery with no date.
Wrangler's directives are written one per line and applied in order. They can be generated from each column's context menu, but it is far quicker to type them:
-- 1) Split the line by the separator and name the columns parse-as-csv :body ';' true drop :body -- 2) Strip stray whitespace from every text column trim :DESTINATARIO trim :CIUDAD trim :INCIDENCIA -- 3) Nulls: the carrier writes '-' where there is no data find-and-replace :FECHA_ENTREGA s/^-$//g find-and-replace :INCIDENCIA s/^-$//g set-column :INCIDENCIA (INCIDENCIA == null || INCIDENCIA.isEmpty()) ? null : INCIDENCIA -- 4) European dates -> a real date type parse-as-simple-date :FECHA_ENTREGA dd/MM/yyyy format-date :FECHA_ENTREGA yyyy-MM-dd -- 5) Decimal commas -> full stops, then to a number find-and-replace :PESO s/,/./g find-and-replace :IMPORTE s/,/./g set-type :PESO double set-type :IMPORTE double -- 6) Split 'SURNAME, FIRST NAME' into two columns split-to-columns :DESTINATARIO , rename :DESTINATARIO_1 apellidos rename :DESTINATARIO_2 nombre trim :apellidos trim :nombre -- 7) Normalise the format of the personal names titlecase :apellidos titlecase :nombre uppercase :CIUDAD -- 8) Derived column: there was an incident if the field is not empty set-column :hubo_incidencia (INCIDENCIA != null) -- 9) Column names in lower case and consistent with the rest of the warehouse 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) Discard rows with no order identifier: they are of no use at all filter-rows-on condition-false pedido_id != null && !pedido_id.isEmpty() -- 11) Drop the first name and surname columns for DATA MINIMISATION drop :apellidos drop :nombre
The most useful directives, grouped:
| Family | Directives | What for |
|---|---|---|
| Parsing | parse-as-csv, parse-as-json, parse-as-xml, parse-as-fixed-length |
Turn raw text into columns |
| Dates | parse-as-simple-date, format-date, parse-as-datetime |
Regional formats |
| Text | trim, uppercase, lowercase, titlecase, cleanse-column-names |
Normalisation |
| Search | find-and-replace, extract-regex-groups, split-to-columns |
Regular expressions |
| Types | set-type, fill-null-or-empty |
Conversion and nulls |
| Filters | filter-rows-on, filter-row-if-matched |
Discard rows |
| Columns | rename, drop, keep, set-column, merge |
Structure |
| Quality | send-to-error |
Divert invalid records |
Steps 10 and 11 of the example deserve a comment.
Step 10 uses filter-rows-on, which discards silently. If you would rather keep what was discarded so you can review it — and normally you should — use send-to-error, which sends the record to the pipeline's Error Collector instead of throwing it away:
Step 11 is a compliance decision, not a cleansing one:
GDPR warning. The carrier's CSV contains the recipient's first name and surname, which is identifying personal data. For analysing delivery times and incidents, that data adds absolutely nothing: the postcode and the order identifier are enough. Removing it at the cleansing stage, before it reaches the analytical warehouse, is data minimisation by design and is the correct practice. Any processing of real personal data must be reviewed by a compliance professional or the DPO before going to production. All the data in this course is fictitious.
Wrangler's advantage over writing this in Python or SQL is not power — Python can do all of this and more — but the feedback loop: you apply a directive and immediately see its effect on 100 real rows. When parse-as-simple-date fails because there is a differently formatted date on row 47, you see it there and then, not once the pipeline has been running in production for twenty minutes.
- The real pipeline: the ERP's MySQL into
alpinashop_analitica
alpinashop_analiticaLet us tackle the integration Lucía really wants: bringing in the purchase costs and the stock from the ERP so she can calculate the real margin per product.
Block 1 — Source: Database. Configured with:
- Plugin type:
Database(generic JDBC). - JDBC driver: the MySQL connector, which has to be uploaded beforehand from the Hub or as a custom plugin.
- Connection string:
jdbc:mysql://10.20.0.15:3306/erp_almacen(private IP reachable over the VPN). - Import Query: the query that extracts the data.
-- ERP extraction query.
-- $CONDITIONS is MANDATORY if partitioned reading is used: Data Fusion
-- replaces it with a different range in each parallel task.
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 $CONDITIONSTwo fundamental details of this configuration:
${fecha_desde}and${fecha_hasta}are Data Fusion macros: runtime arguments. They let the same pipeline serve both the full initial load and the daily incremental one, without duplicating it. The orchestrator in 04-06 will pass it the values.$CONDITIONSwithnumSplits: if you setSplit-By Field Name = skuandNumber of Splits = 4, Data Fusion fires four queries in parallel over different ranges. That speeds up the initial load enormously, but it puts four times as much load on the ERP's MySQL. With a production database that is also serving the warehouse weighing scale, you have to be prudent: one or two splits, and run in the small hours.
Block 2 — Transform: Wrangler. Cleansing specific to the ERP:
-- The ERP mixes upper and lower case in the SKUs uppercase :sku trim :sku -- Internal movement codes -> readable labels set-column :tipo_movimiento (tipo_movimiento == 'E' ? 'entrada' : (tipo_movimiento == 'S' ? 'salida' : 'ajuste')) -- The ERP stores costs in cents as an integer set-column :coste_unitario_eur coste_unitario / 100.0 drop :coste_unitario -- Discard movements with no SKU: they are accounting adjustments, not stock ones send-to-error sku == null || sku.isEmpty() -- Stamp of when it was extracted, so loads can be audited set-column :cargado_en datetime:CurrentDateTime()
Block 3 — Analytics: Group By. Weighted average cost 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
Block 4 — Sink: BigQuery.
- Dataset:
alpinashop_analitica - Table:
costes_producto_erp - Operation:
Upsertkeyed onsku(updates if it exists, inserts if not) - Truncate Table: disabled
- Service Account:
[email protected] - Location:
europe-west1
The Upsert option is what makes the pipeline re-runnable without duplicating. It is the equivalent of SQL's MERGE, and it is what lets you launch the pipeline twice without spoiling anything — the same idempotency property we were chasing in 04-04.
Block 5 — Error Collector → Sink GCS. The records diverted with send-to-error go to gs://alpinashop-datalake/cuarentena/erp/, in JSON format, with the reason for the rejection.
And with that data finally in BigQuery, Lucía can at last answer her question with SQL:
-- Real margin per product: average selling price against the ERP cost
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; -- worst margins first: that is what is actionableThat query, sorted by the lowest margin, is exactly the kind of result that changes decisions: products that sell a lot and leave little behind. It was not possible before this lesson because the cost lived on a server in Sabadell.
- Connectors and plugins from the Hub
The Hub is the catalogue from which plugins are installed into the instance. The main categories:
| Type | Available examples |
|---|---|
| Databases | 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 |
| Files and protocols | HTTP, FTP/SFTP, Amazon S3, Azure Blob, Excel, XML |
| Transformations | Wrangler, JavaScript, Python Evaluator, XML Parser, Validator |
| Analytics | Joiner, Group By, Deduplicate, Pivot, Window Aggregation |
For AlpinaShop's third source, the backpack supplier's API, the HTTP plugin is used:
- URL:
https://api.proveedor-montana.example/v2/catalogo - HTTP Method:
GET - Headers:
Authorization: Bearer ${api_token}— where${api_token}is a macro resolved from Secret Manager (03-06), never written into the pipeline. - Format:
json - JSON/XML Result Path:
$.productos— the path within the response where the array lives. - Pagination Type:
Link HeaderorIncrement an Index, depending on what the API supports.
Pagination is the thing most often forgotten: without configuring it, only the first N results are fetched and nobody notices until products go missing.
If a source has no plugin, there are three ways out: write one in Java (CDAP is extensible), use the Python Evaluator for one-off transformations, or — the most sensible — export the data to Cloud Storage with a script and read it from there. The last one is usually the right answer for an exotic, low-volume source.
- Deploying, running and scheduling
A pipeline in Data Fusion has two states: draft (editable, previewable) and deployed (immutable, executable, versioned).
The workflow:
- Preview. Runs the pipeline over a small sample without writing to the destination. It shows the data coming out of each block. It is the equivalent of Beam's
DirectRunnerand you should always use it before deploying. - Deploy. Freezes the pipeline with a version number. To change it, a new version is created.
- Run. Executes it. Runtime arguments (the macros) can be passed in.
- Schedule. Schedules recurring runs with a cron expression.
All of this can also be done over the REST API, which is how the orchestrator in 04-06 will invoke it:
INSTANCE=$(gcloud data-fusion instances describe alpinashop-fusion \
--location=europe-west1 --format="value(apiEndpoint)")
TOKEN=$(gcloud auth print-access-token)
# Run the pipeline with the date macros
curl -X POST \
-H "Authorization: Bearer ${TOKEN}" \
-H "Content-Type: application/json" \
"${INSTANCE}/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"
}'
# Check the status of the runs
curl -H "Authorization: Bearer ${TOKEN}" \
"${INSTANCE}/v3/namespaces/default/apps/erp-costes-a-bigquery/workflows/DataPipelineWorkflow/runs"Compute profiles define what the Dataproc cluster running the pipeline will look like: number of workers, machine type, network, service account. This is where the execution cost is controlled:
Profile "alpinashop-perfil-computo" Provisioner : Dataproc Region : europe-west1 Master : n2-standard-2, 1 node, 100 GB disk Workers : n2-standard-2, 2 nodes, 100 GB disk Network : alpinashop-vpc / sn-datos-euw1 Internal IP only : true Service Account : [email protected] Image version : 2.2-debian12 Idle TTL : 10 minutes
That profile applies exactly what we learned in 04-03: private network with no public IP, its own service account, pinned image version and self-destruction on idleness.
- What happens underneath: ephemeral Dataproc
When you press Run, this is what actually happens:
sequenceDiagram
participant U as Lucia
participant DF as Data Fusion
participant DP as Dataproc
participant BQ as BigQuery
U->>DF: Run the pipeline
DF->>DF: translates the visual graph into a Spark job
DF->>DP: creates an ephemeral cluster (2-5 min)
DP->>DP: runs the Spark job
DP->>BQ: writes the result
DP-->>DF: job finished
DF->>DP: destroys the cluster (Idle TTL)
DF-->>U: SUCCEEDED status + lineage recorded
Very practical consequences of knowing this:
Start-up takes 2 to 5 minutes. That is not Data Fusion being slow: it is the time to create the cluster. That is why Data Fusion is no use whatsoever for anything that needs low latency. A pipeline that moves 200 rows takes almost as long as one that moves 20 million, because the dominant cost is the start-up.
The real cost is the sum of two things: the instance hours (always) plus the Dataproc hours (per run). A daily 8-minute pipeline with 3 nodes is a few cents of Dataproc, but the instance keeps billing 24 hours a day.
You can reuse a cluster. By configuring a profile that points to an existing cluster instead of creating one, the start-up time disappears. It makes sense when many pipelines run back to back, for instance in a nightly window: create the cluster, launch ten chained pipelines, destroy it.
Spark errors show up in the Dataproc logs. When a pipeline fails on memory (OutOfMemoryError in an executor) or on data skew, the diagnosis is exactly the one from 04-03: look at the Spark interface and the stages. Data Fusion does not hide that reality from you, it only wraps it.
- Field-level data lineage
For many organizations, this is the main reason to pay for Data Fusion.
Lineage answers two questions that sound trivial and that almost no company can answer:
- "Where exactly does this field in the report come from?"
- "If I change this column in the ERP, what breaks?"
Data Fusion records lineage automatically as each pipeline runs, at two levels:
Dataset lineage: which sources feed which destinations, with which pipeline and when.
Field lineage: which specific source column produces which destination column, and what operations it went through along the way.
flowchart LR
A["erp_almacen.movimientos_stock<br/>coste_unitario (INT, cents)"]
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
Nobody draws that diagram: Data Fusion generates it from the run. And it answers the question that in many companies costs days of archaeology: if somebody asks why the margin in the management report looks odd, the lineage shows that coste_medio_eur comes from a division by 100 and an unweighted average — which, incidentally, is a debatable decision that the lineage leaves in plain sight.
Why it matters so much:
| Scenario | Without lineage | With lineage |
|---|---|---|
| Audit | "We think it comes from the ERP" | Documented, dated traceability |
| Change in the source | Deploy and see what breaks | You know in advance which pipelines depend on it |
| Suspicious figure | Days of investigation | One click back to the source |
| GDPR: where a piece of personal data lives | Manual search everywhere | The field is traced across every destination |
| Decommissioning a system | Fear of switching it off | You see exactly what consumes from it |
That last row is especially valuable in a migration: knowing what depends on the ERP before touching it. And the one before it connects directly with what we will see in 04-07, where Dataplex consolidates the lineage of the whole platform, not just Data Fusion's.
- Replication with CDC from the ERP's MySQL
A batch pipeline that runs every night leaves the data up to 24 hours behind, and on top of that loads the ERP with a heavy query. The alternative is change data capture (CDC): instead of querying the table, you read the database's transaction log and replicate the changes as they happen.
flowchart LR
M["MySQL ERP Sabadell<br/>binlog"]
R["Data Fusion Replication<br/>reads the binlog"]
S["Staging table<br/>in BigQuery"]
B["Target table<br/>alpinashop_analitica.stock_erp"]
M -->|INSERT/UPDATE/DELETE| R
R -->|change events| S
S -->|periodic MERGE| B
Preparation on the ERP's MySQL:
-- Requirements on the source server (applied by the ERP administrator)
-- In my.cnf:
-- server-id = 1
-- log_bin = mysql-bin
-- binlog_format = ROW <-- ESSENTIAL: ROW, not 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 is non-negotiable: with STATEMENT, the binlog stores the SQL statements, not the resulting rows, and replication cannot rebuild the state. It is the first requirement to verify with the ERP vendor.
In Data Fusion you use the Replication module (Enterprise edition):
- Source: MySQL, with the connection and the
cdc_datafusionuser. - Table selection:
movimientos_stock,proveedores,articulos. - Destination: BigQuery, dataset
alpinashop_analitica, with a staging prefix. - Source assessment: checks permissions and configuration before starting.
- Initial snapshot + continuous streaming.
Data Fusion first loads a full snapshot of the tables and then applies the changes continuously, writing into staging tables and running a periodic MERGE against the final tables.
Important considerations about CDC, to weigh up before committing:
| Aspect | Reality |
|---|---|
| Latency | Seconds or a few minutes, against 24 hours for the batch |
| Load on the source | Much lower: reading the binlog runs no queries |
| Deletes | They are captured, which an incremental load by date does not do |
| Cost | Requires the Enterprise edition and continuous execution: it is expensive |
| Requirements | Source server configuration; sometimes the vendor will not allow it |
| Schema changes | They are handled, but require attention |
That point about deletes is the strongest technical argument in favour of CDC. An incremental pipeline reading WHERE fecha_modificacion > X never finds out that a row has been deleted, and the analytical warehouse accumulates phantom records indefinitely. CDC does see the DELETE.
The alternative worth considering is Datastream, a Google service dedicated exclusively to CDC, simpler and cheaper than switching on Enterprise in Data Fusion just for this. It is in the table in the next section.
- When Data Fusion, when Dataflow, when Datastream, when
bq load
bq loadThe honest table, which is what you should take away from this lesson:
| Criterion | bq load |
Datastream | Data Fusion | Dataflow |
|---|---|---|---|---|
| Typical case | Ready-made file → table | DB replica with CDC | Integrating varied sources with cleansing | Bespoke transformation, streaming |
| Code | One command | None | None (visual) | Python or Java |
| Profile needed | Anyone | Anyone + source DBA | Analyst | Data engineer |
| Fixed cost | Zero | Per GB processed | Per instance hour | Zero in batch |
| Variable cost | Free | Low | Dataproc per run | Workers per hour |
| Minimum latency | Minutes | Seconds | 2-5 min start-up | Seconds (streaming) |
| Transformations | None | None (replicates as is) | Many, visual | Unlimited |
| Real streaming | No | Yes (replication) | Limited | Yes, its strong point |
| External connectors | No | Databases | Very many | Whichever you program |
| Lineage | No | No | Yes, per field | Via Dataplex |
| Versioning in Git | Trivial | N/A | Awkward (exported JSON) | Trivial |
And translated into decision rules:
Use bq load when the data is already in Cloud Storage in the right shape. It is free and there is nothing to maintain. Always start by asking yourself whether this is enough. Half the pipelines that exist in the world are unnecessary.
Use Datastream when the goal is to replicate a whole database into BigQuery with low latency and delete capture, without transforming anything. It is simpler and cheaper than Data Fusion Enterprise for that specific task, and it supports MySQL, PostgreSQL, Oracle and SQL Server.
Use Data Fusion when there are several heterogeneous sources, non-trivial transformations and cleansing are needed, the team has no programming profile, and lineage is a real requirement (audit, compliance). And when the volume of work justifies the cost of the instance.
Use Dataflow when there is streaming with event-time semantics, complex or state-dependent transformations, or when the team prefers versionable code with automated tests.
And there is one very sensible combination that you see a lot in practice: Datastream to replicate the raw data + BigQuery SQL to transform it (pure ELT). No Data Fusion, no Dataflow, no instances to pay for. For many cases, including a good part of AlpinaShop's, it is the most efficient answer.
- AlpinaShop's reasoned decision
With all of the above on the table, Lucía and Marta decide as follows:
The Sabadell ERP → Datastream, not Data Fusion. The requirement is to replicate three tables with low latency and capture deletes. There is no transformation: the costs are calculated afterwards with SQL in BigQuery, where Lucía is perfectly at home. Datastream costs a fraction of an Enterprise instance and there is no instance to switch off. The transformations from section 8 (upper case, cents to euros, movement type labels) become a BigQuery view, which as a bonus ends up versioned in Git.
The carrier's monthly CSV → Data Fusion, with an ephemeral instance. Here it wins: the file is genuinely dirty and the eleven Wrangler directives are written in twenty minutes while looking at the data, against a couple of days of development and testing in Beam. It is monthly, so the instance is switched on, run and switched off. With the Basic edition and under 120 hours a month, the cost is zero. The instance is managed like this:
# Stop the instance when it is not in use (Enterprise and Basic both allow it)
gcloud data-fusion instances update alpinashop-fusion \
--location=europe-west1 --enable-instance-stopped
# And start it again on the day of the monthly load
gcloud data-fusion instances update alpinashop-fusion \
--location=europe-west1 --no-enable-instance-stoppedThe backpack supplier's API → a scheduled Cloud Function. That is 400 products in JSON, once a week. Spinning up a three-node Dataproc cluster for five minutes to read 400 records is disproportionate. Forty lines of Python in Cloud Functions (06-03) triggered by Cloud Scheduler solve the case for cents.
The general rule AlpinaShop adopts, which works as a reusable criterion:
- Is the data already in Cloud Storage in table shape? →
bq load. - Do you need to replicate a database without transforming it? → Datastream + BigQuery views.
- Is it a small volume from an API? → scheduled Cloud Function.
- Is it a dirty, recurring file that an analyst is going to clean? → Data Fusion with an ephemeral instance.
- Is there streaming, windowing or complex logic? → Dataflow.
- Is it algorithmic with Spark libraries? → Dataproc Serverless.
That list, with the discipline of going through it in order and stopping at the first option that works, is what stops a small business's data platform from ending up costing what a multinational's does.
Common Mistakes and Tips
Leaving the instance switched on. It is the expensive mistake with this service and it bears repeating: a forgotten Basic instance costs over a thousand euros a month without running a single pipeline. Switch it off or delete it.
Choosing Enterprise without needing it. It is only required for CDC, high concurrency and high availability. Basic covers most cases in a small business at less than half the price.
Using Data Fusion for tiny volumes. A Dataproc cluster to process 400 rows is disproportionate. The start-up takes longer than the work.
Not versioning the pipelines. Even though they are visual, they should be exported to JSON and kept in Git. Without that, a deleted instance takes months of work with it and there is no change history.
# Export a pipeline so it can be versioned
curl -H "Authorization: Bearer $(gcloud auth print-access-token)" \
"${INSTANCE}/v3/namespaces/default/apps/erp-costes-a-bigquery" \
> pipelines/erp-costes-a-bigquery.jsonParallelising the extraction without measuring the impact on the source. Number of Splits = 8 against the ERP's MySQL can take the warehouse weighing scale out of service. Start at 1 and go up while measuring.
Writing credentials into the pipeline configuration. Use macros resolved from Secret Manager. A pipeline exported to JSON with the password inside ends up in Git, and that is a security incident.
Forgetting pagination in the HTTP plugin. You fetch the first 100 records and nobody notices until products go missing from a report.
Skipping the Preview. It is free, takes seconds and shows exactly what comes out of each block. Deploying without previewing means launching a cluster to discover a type error.
Not using Upsert in the sink. With Insert, every re-run duplicates rows. With Upsert on the business key, the pipeline is idempotent and can be relaunched without fear.
Tip: use macros from the very start. Dates, paths, table names and credentials as ${variable}. It turns a rigid pipeline into a reusable, orchestrable one.
Tip: always put an Error Collector in. Rejected records must go somewhere, with their reason. It is the same discipline as Beam's quarantine in 04-02.
Tip: review the lineage after every deployment. It is the best way to verify that the pipeline does what you think it does, and it costs nothing.
Exercises
Exercise 1: Wrangler directives for the supplier's file
The backpack supplier sends this tab-separated text file:
REF DESCRIPCION PVP_RECOMENDADO COSTE STOCK ALTA ACTIVO mb-4001 Trekking Backpack 40L Blue 89,90 EUR 52,30 EUR 120 01-03-2024 S MB-4002 trekking backpack 30l red 74,50 EUR 43,10 EUR 0 15-06-2025 S MB-4003 Alpine Backpack 55L 129,00 EUR 78,00 EUR - N/D N
Write the Wrangler directives that produce a clean dataset with: sku in upper case with no whitespace; nombre in title case; pvp_eur and coste_eur as decimal numbers without the word EUR; stock as an integer, with - converted to 0; fecha_alta as a yyyy-MM-dd date, tolerating N/D as null; activo as a boolean; and a calculated column margen_pct. Records with no REF must be diverted to error.
Exercise 2: designing the shipments pipeline
Design the complete pipeline that reads the carrier's monthly CSV from gs://alpinashop-datalake/transportista/2026/03/envios.csv, cleanses it, joins it with the pedidos table in alpinashop_analitica to add the country and the channel, calculates the delay in days between fecha_pedido and fecha_entrega, and writes into alpinashop_analitica.envios. State: each block with its type and essential configuration, how to handle shipments whose pedido_id does not exist in pedidos, which operation to use in the sink and why, and which macros you would define.
Exercise 3: the tool decision, with an economic justification
AlpinaShop absorbs a competitor and inherits four new integrations. For each one, choose between bq load, Datastream, Data Fusion, Dataflow, Dataproc Serverless or Cloud Function, and justify it including an order-of-magnitude estimate of the monthly cost:
- An 80 GB PostgreSQL 14 with the competitor's customer history, which must be replicated into BigQuery with under 5 minutes of lag and capturing deletes.
- A weekly Excel file of 3,000 rows sent by the purchasing team, with columns that get renamed every other week and hand-typed data.
- A stream of 2,000 events per second from a mobile app that has to be aggregated into 5-minute windows before being stored.
- A nightly 12 GB Parquet dump that the competitor already drops into a bucket, with the exact schema of the target table.
Solutions
Solution 1
-- 1) Parse the tab-separated file, with a header parse-as-csv :body '\t' true drop :body -- 2) SKU: strip whitespace and normalise to upper case trim :REF uppercase :REF rename :REF sku -- 3) Divert records with no reference to error send-to-error sku == null || sku.isEmpty() -- 4) Name: clean whitespace and apply title case trim :DESCRIPCION titlecase :DESCRIPCION rename :DESCRIPCION nombre -- 5) Amounts: strip ' EUR', swap the decimal comma and convert to a number 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: '-' means zero, not null find-and-replace :STOCK s/^-$/0/g set-type :STOCK int rename :STOCK stock -- 7) Date: 'N/D' to null, and dd-MM-yyyy format to a real date 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) Active: 'S'/'N' to boolean set-column :ACTIVO (ACTIVO == 'S') rename :ACTIVO activo -- 9) Calculated column: percentage margin, protected against division by zero set-column :margen_pct (pvp_eur != null && pvp_eur > 0) ? ((pvp_eur - coste_eur) / pvp_eur * 100) : null
Expected result over the three sample rows:
| sku | nombre | pvp_eur | coste_eur | stock | fecha_alta | activo | margen_pct |
|---|---|---|---|---|---|---|---|
| MB-4001 | Trekking Backpack 40L Blue | 89.90 | 52.30 | 120 | 2024-03-01 | true | 41.8 |
| MB-4002 | Trekking Backpack 30L Red | 74.50 | 43.10 | 0 | 2025-06-15 | true | 42.1 |
| MB-4003 | Alpine Backpack 55L | 129.00 | 78.00 | 0 | null | false | 39.5 |
The three points being assessed: distinguishing - (which means zero units) from N/D (which means unknown data, that is, null) — confusing them would falsify any stock report; protecting the division in the calculated column; and diverting to error instead of filtering silently, so that a file with empty references leaves a trace.
Solution 2
Block 1 — Source: GCS
- Path:
gs://alpinashop-datalake/transportista/${anyo}/${mes}/envios.csv - Format:
text(one row per line; parsing is done in Wrangler because the separator is;) - Service Account:
sa-fusion-pipelines@...
Block 2 — Transform: Wrangler
The directives from section 7, including the drop of first name and surname for data minimisation, and send-to-error for rows with no pedido_id.
Block 3 — Source: BigQuery
- Dataset/Table:
alpinashop_analitica.pedidos - Import Query (better than reading the whole table):
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}'The partition filter is mandatory: the table has require_partition_filter=TRUE since 04-01, so without it the pipeline would fail. And even if it did not, reading the whole table every month would be throwing money away.
Block 4 — Analytics: Joiner
- Inputs: the Wrangler output (left) and the BigQuery output (right)
- Join type:
Left Outeronpedido_id - Fields: everything from the clean CSV, plus
fecha_pedido,paisandcanalfrom the right-hand side
The join type is the key decision in this exercise. With Inner, shipments whose pedido_id does not exist in pedidos would disappear without a trace, and nobody would know that shipments were missing. With Left Outer they are all kept, with pais and canal null, and those nulls are the signal that there is an integrity problem to investigate: orders from months outside the date range, or identifiers mistyped in the carrier's file.
Block 5 — Transform: Wrangler (second one)
-- Delay in days between order and delivery set-column :retraso_dias (fecha_entrega != null && fecha_pedido != null) ? (dateDiff(fecha_entrega, fecha_pedido)) : null -- Orphan shipment flag, so they can be counted in a quality report set-column :sin_pedido_asociado (pais == null)
Block 6 — Sink: BigQuery
- Table:
alpinashop_analitica.envios - Operation:
Upsertkeyed onenvio_id - Partition field:
fecha_entrega; Cluster field:pais
Why Upsert. The carrier's file tends to be resent corrected when there are errors, and it sometimes includes shipments from the previous month that were delivered late. With Insert, every resend would duplicate rows and the delivery-time report would lie. With Upsert on envio_id, re-running the pipeline as many times as necessary always produces the same result: it is idempotent, exactly the same principle we demanded of the Pub/Sub consumers in 04-04.
Block 7 — Error Collector → Sink GCS
- Path:
gs://alpinashop-datalake/cuarentena/envios/${anyo}/${mes}/ - Format:
json
Macros to define: ${anyo}, ${mes}, ${fecha_desde}, ${fecha_hasta}. With them, the same pipeline serves any month and a whole history can be reprocessed by changing only the arguments. Without them, you would have to duplicate the pipeline or edit it every month.
Solution 3
1. An 80 GB PostgreSQL replicated into BigQuery, under 5 minutes of lag, with deletes → Datastream.
That is literally its definition: managed CDC over PostgreSQL, code-free, with DELETE capture — which an incremental load by date would never detect. Data Fusion Enterprise could also do it, but it would demand the expensive edition (~$3,000/month of instance) to do exactly the same thing. Datastream is billed per GB processed: the initial 80 GB load plus the daily changes put the cost in the order of a few tens of euros a month. The subsequent transformation is done with views in BigQuery, free in effort and versionable in Git.
2. A weekly Excel of 3,000 rows with changing columns and hand-typed data → Data Fusion with an ephemeral instance. This is the textbook case for Wrangler: dirty data typed by people, an unstable schema, and an analyst — not a programmer — who needs to fix it while looking at the data. Programming it in Beam would mean redeploying code every time purchasing renamed a column. With a Basic instance switched on for one hour a week, that is 4 hours a month, well inside the 120 free ones: effective cost zero, plus a few cents of Dataproc per run. The key word in the answer is ephemeral: if the instance is left switched on, the same solution costs over a thousand euros a month.
3. 2,000 events per second aggregated into 5-minute windows → Dataflow.
Streaming with time windows: Beam's exclusive territory. Data Fusion does not do this well, bq load does not apply and a Cloud Function cannot hold window state. A streaming pipeline with 2-3 permanent workers is in the order of €100-150 a month, a cost that has to be taken on consciously because it is the only tool that meets the requirement. If the business could tolerate 15 minutes of lag instead of 5, the alternative — a Pub/Sub subscription to Cloud Storage plus micro-batches — would cost a fraction; it is worth asking before switching the pipeline on.
4. A nightly 12 GB Parquet with the exact schema of the target table → bq load.
With no transformation and the schema already correct, there is nothing to process. Batch loading into BigQuery is free and Parquet carries the schema inside it, so you do not even have to declare it. One command in the 04-06 orchestrator, triggered by the bucket notification we set up in 04-04. Cost: €0 of processing, only the storage for the 12 GB. Any other option on this list would be paying to make a copy, and that is exactly the reflex to correct: the first question in the face of any integration is "is bq load enough?".
Conclusion
You have seen the least showy and most real side of a data platform: the data that is not born in the cloud. The ERP's MySQL in Sabadell, the carrier's CSV with its European dates and its nulls written as a hyphen, and the backpack supplier's API. None of them publishes to Pub/Sub, none of them writes Parquet, and Lucía needs all three to calculate the real margin per product and explain why reviews are dropping in certain areas.
You have understood what a visual ETL/ELT is and, above all, who it serves: profiles who have mastered the business and SQL but are not going to write Beam, and organizations that need documented traceability. And also who it does not serve, which matters just as much. You know that Data Fusion is managed CDAP, with its Studio, its Wrangler, its Hub and its lineage record, and that underneath it executes nothing: it translates the visual graph into Spark and spins up an ephemeral Dataproc, which explains its 2-5 minutes of start-up and why it is no use for low latencies.
You have put the cost up front instead of hiding it: Developer, Basic and Enterprise editions, billed per hour that the instance exists, used or not, with a permanent Basic instance costing more than the whole of the rest of AlpinaShop's platform put together. That fact does not disqualify the product — for a company with twenty sources it is cheap against the salaries it saves — but it does force the ephemeral instance strategy in a small business.
You have cleansed the carrier's CSV with Wrangler and its directives: semicolon parsing, trim, nulls written as a hyphen, dd/MM/yyyy dates converted to a date type, decimal commas, splitting the recipient column, title case, derived columns, and the final drop of first name and surname for data minimisation, with the explicit GDPR warning and the requirement for compliance review. You have weighed up the tool's real advantage, which is not power but the feedback loop: you see the effect of each directive on real data instantly.
You have designed the ERP pipeline with its date macros, its partitioned reading with the prudence not to knock over the warehouse weighing scale, its cleansing, its aggregation and its Upsert sink so that it is idempotent and re-runnable. You know the Hub and its connectors, the HTTP plugin with its easily forgotten pagination, the Preview → Deploy → Run → Schedule flow, the compute profiles that control the underlying cluster, and the REST API through which the orchestrator will invoke it. You have seen field-level lineage, which answers "where does this column come from?" and "what breaks if I change this?" without days of archaeology, and CDC replication from MySQL with its non-negotiable binlog_format = ROW and its great argument: it is the only way to find out about deletes.
And you have closed with the honest table and with AlpinaShop's reasoned decision, which is not "let us use Data Fusion for everything" but an ordered list: first bq load if it is enough, then Datastream if it is replication, then a Cloud Function if it is a small job, then Data Fusion with an ephemeral instance if the file is dirty and an analyst is cleaning it, and Dataflow or Dataproc if there is streaming or algorithms. Going through that list in order and stopping at the first option that works is what separates a proportionate data platform from a ruinously expensive one.
And now a new problem appears, one that is a direct consequence of the success of the previous four lessons. AlpinaShop has a nightly Cloud SQL export, a load into BigQuery, a Dataflow pipeline, a Spark job on Dataproc Serverless, a monthly Data Fusion pipeline, several aggregation queries and a materialised view to refresh. That is seven or eight processes that depend on one another: there is no sense in launching the aggregation before the load has finished, nor in refreshing the business view with half-loaded data.
Right now that is handled by a cron on a virtual machine that launches scripts at fixed times, worked out by eye with plenty of margin. It works until the first night the export takes twenty minutes longer than usual: then the load reads an incomplete file, the aggregation computes over partial data, the management report wakes up with false figures, and nobody finds out until somebody looks at them mid-morning.
In 04-06, Cloud Composer and Workflows, we will solve that. You will see why a cron is not an orchestrator and what coordinating dependencies, retries, alerts and reprocessing really means. You will meet Cloud Composer — managed Apache Airflow — with its DAGs, tasks, operators and sensors, and you will write the complete DAG for AlpinaShop's nightly process commented line by line, with the clear warning that Composer is expensive for a small business. And you will meet Workflows, serverless orchestration in YAML with no fixed cost, with Cloud Scheduler for the trigger, which is where AlpinaShop is going to start.
Google Cloud Platform (GCP) Course
Module 1: Introduction to Google Cloud Platform
- What is Google Cloud Platform?
- Setting Up Your GCP Account
- A Tour of the GCP Console
- Projects, Resource Hierarchy and Billing
- Regions, Zones and the Shared Responsibility Model
- Cloud Shell and the gcloud CLI
Module 2: Core GCP Services
- Compute Engine: Virtual Machines on Google Cloud
- Cloud Storage: Object Storage
- Cloud SQL: Managed Relational Databases
- App Engine: Platform as a Service
- Google Kubernetes Engine (GKE)
- NoSQL Databases: Firestore, Bigtable and Spanner
- How to Choose the Right Compute Service
Module 3: Networking and Security
- VPC Networks
- Cloud Load Balancing
- Cloud CDN
- Identity and Access Management (IAM)
- Cloud Armor
- Secrets and Encryption: Secret Manager and Cloud KMS
- Cloud DNS, TLS Certificates and Publishing Services Securely
Module 4: Data and Analytics
- BigQuery: The Analytical Data Warehouse
- Cloud Dataflow: Batch and Streaming Data Processing
- Cloud Dataproc: Managed Spark and Hadoop
- Cloud Pub/Sub: Asynchronous Messaging
- Cloud Data Fusion: Code-Free Data Integration
- Orchestrating Pipelines with Cloud Composer and Workflows
- Data Governance and Dashboards with Dataplex and Looker Studio
Module 5: Machine Learning and AI
- Vertex AI: The Machine Learning Platform on GCP
- AutoML: Custom Models Without Writing Code
- TensorFlow on GCP: Training and Serving Models
- Natural Language API
- Vision API
- Generative AI on Vertex AI: Gemini Models and Embeddings
- MLOps: From Model to Product with Vertex AI Pipelines
Module 6: DevOps and Monitoring
- Cloud Build: Continuous Integration on GCP
- Cloud Source Repositories and Source Code Management
- Cloud Functions: Serverless Functions
- Cloud Monitoring (formerly Stackdriver): Metrics, Dashboards and Alerts
- Cloud Deployment Manager and Native Infrastructure as Code
- Cloud Logging and Cloud Trace: Logs, Traces and Diagnostics
- Terraform on GCP: Infrastructure as Code in Practice
Module 7: Advanced GCP Topics
- Hybrid and Multicloud with Anthos
- Serverless Computing with Cloud Run
- Advanced Networking: Shared VPC, Peering and Hybrid Connectivity
- Security Best Practices
- Cost Management and Optimization
- Reliability: SLOs, High Availability and Disaster Recovery
- Governance at Scale: Organization, Policies and Auditing
