Lucía has a question that SQL does not answer well: which products are bought together. If a customer puts a pair of crampons in the shopping cart, how likely are they also to buy an ice axe? With that matrix, AlpinaShop could suggest add-ons on the product page, group items into bundles and lay the categories out better. It is a calculation over every possible pair within each order, repeated across two years of history, and it is the kind of problem that starts out elegant in SQL and ends up as a JOIN of a table with itself that nobody wants to maintain.
It is also, as it happens, exactly the kind of problem Spark was invented for: iterative, algorithmic computation over large volumes, expressed in code rather than in queries.
Here it is worth being honest about why this lesson exists. If AlpinaShop were starting from scratch today, it would probably solve this with BigQuery and Dataflow and never touch Spark. But the real world does not start from scratch: there are tens of thousands of companies with Hadoop clusters in their data centres, with years of Spark code in production, with teams that know PySpark and do not know Beam, and with libraries — MLlib, GraphX, the whole scientific Python ecosystem — that have no direct equivalent. Cloud Dataproc is Google's answer to that reality: managed Hadoop and Spark, with the peculiarity that it turns the very idea of a cluster on its head.
In this lesson you will see what Hadoop and Spark are in essence, you will create the alpinashop-spark cluster, you will understand the pattern that makes Dataproc something other than rented Hadoop — the ephemeral cluster over storage in Cloud Storage — you will run the PySpark job that answers Lucía's question, and you will compare Dataproc, Dataproc Serverless and Dataflow with proper judgement.
Contents
- Hadoop and Spark in one page
- Why they still matter in 2026
- What Dataproc adds
- Anatomy of a Dataproc cluster
- Creating
alpinashop-spark - The key pattern: an ephemeral cluster over Cloud Storage
- Submitting jobs:
gcloud dataproc jobs submit - AlpinaShop's PySpark job: average basket and products bought together
- Initialisation actions and image versions
- Dataproc Serverless for Spark
- Dataproc, Serverless and Dataflow: the decision table
- Notebooks and Spark SQL over BigQuery
- Migrating an on-premises Hadoop to Google Cloud
- Cost and Spot VMs
- Hadoop and Spark in one page
Hadoop was born in 2006 to solve a concrete problem: processing more data than fitted on one machine, using many cheap computers that fail often. It has three pieces:
- HDFS, a distributed file system that splits files into blocks and replicates them (three times by default) across the disks of the cluster's machines.
- YARN, the resource manager that decides which process runs on which machine.
- MapReduce, the original programming model: you split the work into a map phase (transform each record) and a reduce phase (aggregate by key), and the framework handles the distribution and the failures.
MapReduce worked and it was extremely slow, for a design reason: it wrote to disk between every phase. An iterative algorithm that needs twenty passes over the same data did twenty rounds of writing and reading from disk.
Spark appeared in 2014 with the obvious correction: keeping the intermediate data in memory. For the same iterative algorithm, the improvement was one or two orders of magnitude. And it added a far more pleasant API.
Spark today is organised around:
| Component | What it is | Use |
|---|---|---|
| Spark Core / RDD | The original abstraction: a resilient distributed collection | Fine-grained control, operations not expressible in SQL |
| DataFrame / Spark SQL | Tables with a schema and an optimiser (Catalyst) | 90 % of current usage; SQL over distributed data |
| MLlib | Distributed machine learning library | Clustering, recommendation, classification at scale |
| Structured Streaming | Continuous processing with the DataFrame API | The alternative to Beam inside the Spark world |
| GraphX / GraphFrames | Graph algorithms | Networks, paths, communities |
The execution architecture, which you need to have in your head to understand what follows:
flowchart TD
D["Driver<br/>your PySpark program<br/>builds the plan"]
CM["Resource manager (YARN)"]
E1["Executor 1<br/>tasks + in-memory cache"]
E2["Executor 2"]
E3["Executor N"]
S["Storage<br/>HDFS or Cloud Storage"]
D -->|requests resources| CM
CM -->|allocates| E1 & E2 & E3
D -->|sends tasks| E1 & E2 & E3
E1 & E2 & E3 <--> S
The driver runs your code, builds a graph of operations and chops it into tasks. The executors run those tasks over partitions of the data. Just as in Beam, transformations are lazy: nothing happens until you call an action (count(), collect(), write()). That laziness is what lets the optimiser reorder and fuse operations.
- Why they still matter in 2026
With BigQuery and Dataflow available, why learn this? Four concrete reasons and one consequence.
Existing code. A company migrating to the cloud with 200,000 lines of PySpark in production is not going to rewrite them. Rewriting adds no business value, introduces bugs and eats months. Dataproc lets you move those workloads without touching the code, and optimisation can come later.
People. The market has many more engineers who know Spark than engineers who know Beam. If AlpinaShop's data team hires tomorrow, the candidate is more likely to bring Spark with them. Choosing the technology your team knows how to use is a legitimate architectural decision, not a concession.
MLlib and the Python ecosystem. Algorithms such as ALS for recommendation, k-means at scale or FP-Growth for association rules are implemented, tested and distributed. Writing them in Beam would be absurd. And inside a Spark job you can use pandas, NumPy or scikit-learn over specific partitions.
Open formats and ecosystem. Hive, Presto/Trino, HBase, Kafka, Iceberg, Delta Lake, Hudi: a whole world of open tools that speaks Hadoop's language. If AlpinaShop wanted to leave Google Cloud, a data lake in Parquet over object storage processed with Spark travels to any provider unchanged.
The consequence: Dataproc is not a second-class service or a relic. It is the least traumatic migration route and the right tool when the problem is algorithmic and the team knows Spark.
- What Dataproc adds
Building a Hadoop cluster by hand — installing, configuring YARN, sizing HDFS, tuning executor memory, integrating authentication — is weeks of work and a permanent source of maintenance. Dataproc reduces it to one command.
| Aspect | Self-managed Hadoop | Dataproc |
|---|---|---|
| Creation time | Days or weeks | Under 2 minutes |
| Configuration | Manual, component by component | Preconfigured and coherent |
| Scaling | Buy and rack hardware | Change a number; autoscaling available |
| Upgrades | A project in its own right | Managed image versions |
| Cost at rest | The hardware, always | Zero if the cluster does not exist |
| Cloud integration | You have to build it | Cloud Storage, BigQuery, Logging, IAM as standard |
The two minutes of creation time are not a marketing figure: they are what changes the mental model. If creating a cluster takes weeks, the cluster is a permanent installation that has to be looked after, shared between teams and kept switched on at all times just in case. If it takes ninety seconds, the cluster becomes disposable: it is created for one job, destroyed when it finishes, and each job can have its own with the version and the libraries it needs.
That is what we will see in section 6, and it is the most important idea in the lesson.
- Anatomy of a Dataproc cluster
A Dataproc cluster has three types of node:
| Node | Function | Quantity | Notes |
|---|---|---|---|
| Master | Runs the driver, YARN ResourceManager, HDFS NameNode | 1 (or 3 in high availability) | If it goes down with 1 node, the job dies |
| Primary workers | Run tasks; contribute disk to HDFS | Minimum 2 (or 0 in single-node mode) | Standard VMs, stable |
| Secondary workers | Compute only; they do not store HDFS | 0 to N | They can be Spot: up to ~80 % cheaper |
The secondary workers deserve attention. Since they do not take part in HDFS, they can disappear without any data being lost. That is why they can be Spot VMs (the same ones from 02-01, which Google can reclaim with 30 seconds' notice). If one disappears mid-task, YARN reassigns that task to another node and the job carries on, slower but correct.
That is the winning combination for AlpinaShop: a few standard primary workers to provide stability, and many Spot secondaries to provide cheap power.
flowchart TD
subgraph C["Cluster alpinashop-spark"]
M["Master<br/>n2-standard-4<br/>driver + YARN RM"]
W1["Primary worker 1<br/>n2-standard-4"]
W2["Primary worker 2<br/>n2-standard-4"]
S1["Spot secondary 1"]
S2["Spot secondary 2"]
S3["Spot secondary N<br/>autoscaled"]
end
GCS["Cloud Storage<br/>gs://alpinashop-catalogo<br/>gs://alpinashop-datalake"]
BQ["BigQuery<br/>alpinashop_analitica"]
M --- W1 & W2 & S1 & S2 & S3
W1 & W2 & S1 & S2 & S3 <--> GCS
W1 & W2 <--> BQ
Notice that the storage is outside the cluster. That is the next section.
- Creating
alpinashop-spark
alpinashop-sparkgcloud config set project alpinashop-datos
# A dedicated bucket for the data lake and the Spark artefacts
gcloud storage buckets create gs://alpinashop-datalake \
--project=alpinashop-datos --location=europe-west1 \
--uniform-bucket-level-access
# Dedicated service account, with least privilege
gcloud iam service-accounts create sa-dataproc-analitica \
--display-name="Analytics Dataproc clusters"
SA="[email protected]"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/dataproc.worker"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/bigquery.dataEditor"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/bigquery.jobUser"
gcloud storage buckets add-iam-policy-binding gs://alpinashop-datalake \
--member="serviceAccount:$SA" --role="roles/storage.objectAdmin"
gcloud storage buckets add-iam-policy-binding gs://alpinashop-catalogo \
--member="serviceAccount:$SA" --role="roles/storage.objectViewer"And the cluster:
gcloud dataproc clusters create alpinashop-spark \
--region=europe-west1 \
--zone=europe-west1-b \
--service-account="$SA" \
--subnet=sn-datos-euw1 \
--no-address \
--master-machine-type=n2-standard-4 \
--master-boot-disk-size=100GB \
--master-boot-disk-type=pd-balanced \
--num-workers=2 \
--worker-machine-type=n2-standard-4 \
--worker-boot-disk-size=200GB \
--num-secondary-workers=2 \
--secondary-worker-type=spot \
--image-version=2.2-debian12 \
--optional-components=JUPYTER \
--enable-component-gateway \
--bucket=alpinashop-datalake \
--max-idle=30m \
--properties="spark:spark.sql.adaptive.enabled=true,\
spark:spark.dynamicAllocation.enabled=true,\
spark:spark.sql.sources.partitionOverwriteMode=dynamic" \
--labels=entorno=produccion,equipo=datos,centro-coste=analitica,aplicacion=analiticaA review of what matters:
| Option | Why it is there |
|---|---|
--subnet=sn-datos-euw1 + --no-address |
The cluster lives in alpinashop-vpc, with no public IP. It reaches the internet through the Cloud NAT from 03-01 and reaches the APIs through Private Google Access |
--service-account |
Its own identity with minimum permissions, not the default Compute account |
--num-secondary-workers=2 --secondary-worker-type=spot |
Cheap power that can disappear without breaking anything |
--image-version=2.2-debian12 |
Pinned version. Without this, Google would pick the most recent one and a job that worked yesterday could fail tomorrow |
--optional-components=JUPYTER + --enable-component-gateway |
Notebooks reachable from the console with IAM authentication, without opening ports |
--bucket=alpinashop-datalake |
Working bucket for the cluster's logs and temporary files |
--max-idle=30m |
The cluster self-destructs after 30 minutes with no jobs. The most profitable option in the whole command |
spark.sql.adaptive.enabled |
Adaptive execution: Spark readjusts partitions and JOIN strategies at run time. It mitigates data skew considerably |
spark.dynamicAllocation.enabled |
Spark requests and releases executors as it needs them |
Adding autoscaling to the cluster requires a separate policy:
# politica-autoescalado.yaml
workerConfig:
minInstances: 2
maxInstances: 2 # primaries do NOT scale: they provide HDFS
secondaryWorkerConfig:
minInstances: 0
maxInstances: 20 # Spot ones do scale, up to 20
basicAlgorithm:
cooldownPeriod: 2m
yarnConfig:
scaleUpFactor: 1.0 # adds 100 % of what YARN asks for
scaleDownFactor: 0.5 # removes half of the surplus: prudent
gracefulDecommissionTimeout: 10m # waits for in-flight tasks to finishgcloud dataproc autoscaling-policies import pol-autoescalado-analitica \
--region=europe-west1 --source=politica-autoescalado.yaml
gcloud dataproc clusters update alpinashop-spark --region=europe-west1 \
--autoscaling-policy=pol-autoescalado-analiticaThe gracefulDecommissionTimeout matters: without it, shrinking the cluster kills nodes with tasks in flight and those tasks have to be redone. With it, it waits for them to finish. And scaleDownFactor: 0.5 avoids the accordion effect of constantly going up and down.
Verification:
gcloud dataproc clusters list --region=europe-west1 \
--format="table(clusterName, status.state, config.workerConfig.numInstances)"
gcloud dataproc clusters describe alpinashop-spark --region=europe-west1
- The key pattern: an ephemeral cluster over Cloud Storage
Here is the idea that changes everything.
In a traditional Hadoop, storage and compute live on the same machines. HDFS lives on the workers' disks. That had an excellent reason in 2006: moving data over the network was hugely expensive compared with reading it from the local disk, so the principle was "take the compute to the data".
But it has a devastating consequence: if you switch the cluster off, you lose the data. That is why Hadoop clusters are always switched on, even if they only work three hours a day. You pay for the hardware all twenty-four.
In Google Cloud, the internal network is so fast that the premise no longer holds: reading from Cloud Storage is not significantly slower than reading from a local disk. That allows the design to be inverted:
flowchart LR
subgraph Antes["Traditional Hadoop"]
H["Permanent cluster<br/>compute + HDFS<br/>on 24x7"]
end
subgraph Ahora["Dataproc pattern"]
GCS["Cloud Storage<br/>gs://alpinashop-datalake<br/>PERMANENT, cheap"]
C1["Ephemeral cluster A<br/>Spark 3.5<br/>lives 20 min"]
C2["Ephemeral cluster B<br/>Spark 3.3 + library X<br/>lives 5 min"]
end
GCS <--> C1
GCS <--> C2
The Cloud Storage connector comes preinstalled in Dataproc and makes gs:// paths behave like HDFS paths for any Spark or Hadoop code:
# The same code that used to read from HDFS...
df = spark.read.parquet("hdfs:///datos/pedidos/")
# ...reads from Cloud Storage by changing the prefix. Nothing else.
df = spark.read.parquet("gs://alpinashop-datalake/pedidos/")The advantages of separating storage and compute are concrete:
- You pay for compute only when you compute. A cluster that lives 20 minutes a day costs 1.4 % of a permanent one.
- The data outlives the cluster. You can destroy it with complete peace of mind.
- Several clusters over the same data. Lucía's with Spark 3.5 and an external supplier's with an old version, simultaneously, without interfering with each other.
- Far higher durability. Cloud Storage replicates with eleven-nines guarantees; HDFS with three copies on three disks in the same rack does not come close.
- The data is reachable from outside Spark. BigQuery reads it as an external table, Dataflow processes it, the application downloads it.
- No upgrade ceremony. To move to a new Spark version, you create a new cluster. You migrate nothing.
The drawbacks, to be fair: somewhat higher latency on operations involving many small files, and the fact that Cloud Storage has no atomic directory rename, which affects certain write patterns. It is mitigated with columnar formats and reasonably sized files (128-512 MB), which is what you should be doing anyway.
AlpinaShop's rule: HDFS only as temporary working space inside a job. Nothing that has to outlive the cluster is written to HDFS. And --max-idle on every cluster, no exceptions. A cluster forgotten over a weekend costs more than all the rest of the month's analytics.
For scheduled jobs, you do not even keep a cluster: it is created, run and destroyed in a single step with workflow templates:
# 1) Workflow template
gcloud dataproc workflow-templates create wf-cesta-media --region=europe-west1
# 2) Managed cluster: it is born and dies with the workflow
gcloud dataproc workflow-templates set-managed-cluster wf-cesta-media \
--region=europe-west1 \
--cluster-name=cluster-efimero-cesta \
--service-account="$SA" --subnet=sn-datos-euw1 --no-address \
--master-machine-type=n2-standard-4 \
--worker-machine-type=n2-standard-4 --num-workers=2 \
--num-secondary-workers=4 --secondary-worker-type=spot \
--image-version=2.2-debian12
# 3) The job that will run
gcloud dataproc workflow-templates add-job pyspark \
gs://alpinashop-datalake/jobs/cesta_media.py \
--step-id=cesta-media --workflow-template=wf-cesta-media \
--region=europe-west1 \
-- --date-from=2025-01-01 --date-to=2026-03-31
# 4) Run it: creates the cluster, launches the job, destroys the cluster
gcloud dataproc workflow-templates instantiate wf-cesta-media --region=europe-west1That final command is the one the orchestrator in 04-06 will invoke. Total cost: the minutes the job lasts. Zero for the rest of the month.
- Submitting jobs:
gcloud dataproc jobs submit
gcloud dataproc jobs submitDataproc accepts several job types:
# PySpark
gcloud dataproc jobs submit pyspark gs://alpinashop-datalake/jobs/mi_job.py \
--cluster=alpinashop-spark --region=europe-west1 \
-- arg1 arg2
# Spark (Scala or Java JAR)
gcloud dataproc jobs submit spark --cluster=alpinashop-spark --region=europe-west1 \
--class=com.alpinashop.Informe --jars=gs://alpinashop-datalake/jars/informes.jar
# Spark SQL from a file
gcloud dataproc jobs submit spark-sql --cluster=alpinashop-spark \
--region=europe-west1 --file=gs://alpinashop-datalake/sql/ventas.sql
# Hive, Pig and Presto/Trino are available tooEverything after -- is an argument for your program, not for gcloud. It is a very frequent confusion.
Useful options when submitting:
gcloud dataproc jobs submit pyspark gs://alpinashop-datalake/jobs/cesta_media.py \
--cluster=alpinashop-spark \
--region=europe-west1 \
--py-files=gs://alpinashop-datalake/jobs/utilidades.zip \
--files=gs://alpinashop-datalake/config/categorias.json \
--jars=gs://spark-lib/bigquery/spark-3.5-bigquery-0.42.0.jar \
--properties="spark.executor.memory=6g,spark.executor.cores=2,spark.sql.shuffle.partitions=200" \
--labels=proceso=cesta-media \
-- --date-from=2025-01-01 --date-to=2026-03-31 \
--output=gs://alpinashop-datalake/resultados/cesta/--py-files: your own Python modules, packaged as.zipor.egg, distributed to every executor.--files: data files the job needs, reachable by name in the working directory.--jars: Java dependencies, such as the BigQuery connector.spark.sql.shuffle.partitions: the number of partitions after a shuffle. The default (200) is a poor fit for almost everybody: too many for small data, not enough for large. A reasonable rule of thumb is 2-3 times the cluster's total number of cores.
Follow-up:
gcloud dataproc jobs list --region=europe-west1 --cluster=alpinashop-spark \
--format="table(reference.jobId, status.state, statusHistory[0].stateStartTime)"
gcloud dataproc jobs wait JOB_ID --region=europe-west1 # follows the logs liveThe logs go automatically to Cloud Logging (06-06) and the Spark History Server interface stays reachable through the component gateway, even after destroying the cluster if a persistent history server is configured.
- AlpinaShop's PySpark job: average basket and products bought together
Now the real job. It reads the order lines from the lake, computes the average basket per month and country, and builds the product co-occurrence matrix.
"""
cesta_media.py -- AlpinaShop basket analysis with PySpark.
Input : gs://alpinashop-datalake/pedidos/ (Parquet, partitioned by date)
Outputs : gs://alpinashop-datalake/resultados/cesta/ (Parquet)
alpinashop-datos.alpinashop_analitica.productos_juntos (BigQuery)
Submission:
gcloud dataproc jobs submit pyspark gs://.../cesta_media.py \
--cluster=alpinashop-spark --region=europe-west1 \
-- --date-from=2025-01-01 --date-to=2026-03-31
"""
import argparse
from pyspark.sql import SparkSession, functions as F, Window
PROJECT = "alpinashop-datos"
DATASET = "alpinashop_analitica"
LAKE = "gs://alpinashop-datalake"
def create_session():
"""The SparkSession is the entry point. On Dataproc, the resource
configuration and the master come from YARN: there is no need to state them."""
return (
SparkSession.builder
.appName("alpinashop-cesta-media")
# Temporary bucket the BigQuery connector needs in order to write
.config("temporaryGcsBucket", "alpinashop-datalake")
.getOrCreate()
)def load_lines(spark, date_from, date_to):
"""Reads the order lines from the lake and filters them by date.
Parquet stores the schema and per-block statistics, so Spark
can skip whole files: this is 'predicate pushdown'.
"""
lines = (
spark.read.parquet(f"{LAKE}/pedidos/lineas/")
.filter((F.col("fecha_pedido") >= date_from) & (F.col("fecha_pedido") <= date_to))
# Selecting columns early reduces memory and shuffle
.select("pedido_id", "fecha_pedido", "sku", "cantidad",
"precio_unitario", "importe_linea")
)
headers = (
spark.read.parquet(f"{LAKE}/pedidos/cabeceras/")
.filter((F.col("fecha_pedido") >= date_from) & (F.col("fecha_pedido") <= date_to))
.filter(~F.col("estado").isin("cancelado", "devuelto"))
.select("pedido_id", "pais", "canal", "cliente_id")
)
# broadcast(): the period's headers fit in each executor's memory.
# It avoids the JOIN shuffle, just like Beam's side input in 04-02.
return lines.join(F.broadcast(headers), on="pedido_id", how="inner")
def compute_average_basket(df):
"""Average basket per month and country: total amount and number of items."""
per_order = (
df.groupBy("pedido_id", "pais", F.trunc("fecha_pedido", "month").alias("mes"))
.agg(
F.sum("importe_linea").alias("importe_pedido"),
F.sum("cantidad").alias("articulos_pedido"),
F.countDistinct("sku").alias("skus_distintos"),
)
)
return (
per_order.groupBy("mes", "pais")
.agg(
F.count("*").alias("num_pedidos"),
F.round(F.avg("importe_pedido"), 2).alias("cesta_media_eur"),
F.round(F.expr("percentile_approx(importe_pedido, 0.5)"), 2)
.alias("cesta_mediana_eur"),
F.round(F.avg("articulos_pedido"), 2).alias("articulos_medios"),
F.round(F.avg("skus_distintos"), 2).alias("skus_medios"),
)
.orderBy("mes", "pais")
)F.broadcast() deserves attention: it tells Spark to replicate the small DataFrame on every executor instead of shuffling both sides over the network. It is the same optimisation as Beam's side input in 04-02, and it is the difference between a JOIN of seconds and one of minutes. It only works if the small side fits in memory (a few hundred MB at most).
def compute_products_bought_together(df, min_support=20):
"""Co-occurrence matrix: which pairs of SKUs appear in the same order.
The algorithm is a self-join of the set of orders with itself,
with two fundamental performance precautions.
"""
# 1) An order can have the same SKU on several lines: we keep
# unique (order, sku) pairs so as not to count them twice.
order_sku = df.select("pedido_id", "sku").distinct()
# 2) CRITICAL PRECAUTION: discard orders with too many lines.
# An order with 200 SKUs generates 200*199/2 = 19,900 pairs on its own,
# and those rare corporate orders would dominate the calculation and the memory.
sizes = order_sku.groupBy("pedido_id").agg(F.count("*").alias("n_skus"))
valid_orders = sizes.filter((F.col("n_skus") >= 2) & (F.col("n_skus") <= 30))
base = order_sku.join(F.broadcast(valid_orders.select("pedido_id")),
on="pedido_id", how="inner")
left = base.withColumnRenamed("sku", "sku_a")
right = base.withColumnRenamed("sku", "sku_b")
pairs = (
left.join(right, on="pedido_id")
# 3) sku_a < sku_b removes the self-pairs AND the reversed
# duplicates: (A,B) is counted once, not twice as (A,B) and (B,A).
.filter(F.col("sku_a") < F.col("sku_b"))
.groupBy("sku_a", "sku_b")
.agg(F.count("*").alias("veces_juntos"))
.filter(F.col("veces_juntos") >= min_support)
)
# 4) Association rule metrics: support, confidence and lift.
sku_counts = (base.groupBy("sku").agg(F.count("*").alias("veces_total")))
total_orders = base.select("pedido_id").distinct().count()
result = (
pairs
.join(F.broadcast(sku_counts.withColumnRenamed("sku", "sku_a")
.withColumnRenamed("veces_total", "total_a")),
on="sku_a")
.join(F.broadcast(sku_counts.withColumnRenamed("sku", "sku_b")
.withColumnRenamed("veces_total", "total_b")),
on="sku_b")
.withColumn("soporte", F.col("veces_juntos") / F.lit(total_orders))
.withColumn("confianza_a_b", F.col("veces_juntos") / F.col("total_a"))
.withColumn("confianza_b_a", F.col("veces_juntos") / F.col("total_b"))
.withColumn(
"lift",
(F.col("veces_juntos") * F.lit(total_orders))
/ (F.col("total_a") * F.col("total_b")),
)
)
# 5) Top 5 companions for each product, with a window function
window_spec = Window.partitionBy("sku_a").orderBy(F.desc("lift"))
return (
result
.withColumn("puesto", F.row_number().over(window_spec))
.filter(F.col("puesto") <= 5)
.select("sku_a", "sku_b", "veces_juntos",
F.round("soporte", 5).alias("soporte"),
F.round("confianza_a_b", 4).alias("confianza"),
F.round("lift", 3).alias("lift"),
"puesto")
)How the three metrics are interpreted, because these are the ones Lucía will take to the meeting:
| Metric | What it means | AlpinaShop example |
|---|---|---|
| Support | Proportion of orders containing both products | 0.012 → 1.2 % of orders carry crampons and an ice axe |
| Confidence | If they buy A, the probability they buy B | 0.34 → a third of those who buy crampons buy an ice axe |
| Lift | How much more likely they are to go together than by chance | 8.5 → eight and a half times more than expected: an extremely strong association |
Lift is the one to look at. Confidence misleads with bestsellers: if 60 % of orders carry technical socks, any product will show high confidence towards them without any real relationship existing. Lift corrects for each product's popularity. Lift greater than 1 indicates a real association; lift close to 1, independence.
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--date-from", required=True)
parser.add_argument("--date-to", required=True)
parser.add_argument("--output", default=f"{LAKE}/resultados/cesta/")
args = parser.parse_args()
spark = create_session()
spark.sparkContext.setLogLevel("WARN")
df = load_lines(spark, args.date_from, args.date_to)
# cache(): the DataFrame is used in TWO different calculations. Without cache,
# Spark would reread and refilter the whole lake twice.
df.cache()
basket = compute_average_basket(df)
(basket.coalesce(1) # a single file: the result is small
.write.mode("overwrite")
.parquet(f"{args.output}/cesta_media"))
together = compute_products_bought_together(df)
(together.write.format("bigquery")
.option("table", f"{PROJECT}.{DATASET}.productos_juntos")
.option("writeMethod", "direct")
.mode("overwrite")
.save())
df.unpersist()
print(f"Average basket: {basket.count()} rows")
print(f"Product pairs: {together.count()} rows")
spark.stop()
if __name__ == "__main__":
main()Three performance details that have to be internalised:
cache()materialises the DataFrame in the executors' memory. Without it, since Spark is lazy, every later action would recompute the whole chain from the read of the lake. With two consumers, it saves half the work. Andunpersist()at the end, to free the memory.coalesce(1)reduces to a single partition before writing. It is correct only for small results; with large data, concentrating it all in one executor would bring it down. For large volumes you userepartition(n).writeMethod=directuses BigQuery's Storage Write API instead of going through temporary files in the bucket. It is faster and it avoids managing cleanup.
And the corresponding warning: this analysis uses cliente_id and order data. If identifiable personal data is incorporated at any point, the processing must comply with the GDPR and be reviewed by a compliance professional. For basket analysis you do not need to know who the customer is, only what was in the order; keeping it that way is data minimisation by design. All the data in this course is fictitious.
- Initialisation actions and image versions
An initialisation action is a script that runs on every node when the cluster is created. It is used to install dependencies that do not come in the image.
# Your own script in the bucket
cat > init-alpinashop.sh <<'EOF'
#!/bin/bash
set -euxo pipefail
# Python libraries for the basket analysis
pip install --no-cache-dir mlxtend==0.23.1 pyarrow==16.1.0
# On the master only: diagnostic utilities
ROLE=$(/usr/share/google/get_metadata_value attributes/dataproc-role)
if [[ "$ROLE" == "Master" ]]; then
pip install --no-cache-dir jupyterlab-git
fi
EOF
gcloud storage cp init-alpinashop.sh gs://alpinashop-datalake/init/
gcloud dataproc clusters create alpinashop-spark-ml \
--region=europe-west1 --subnet=sn-datos-euw1 --no-address \
--service-account="$SA" \
--initialization-actions=gs://alpinashop-datalake/init/init-alpinashop.sh \
--initialization-action-timeout=10m \
--image-version=2.2-debian12 \
--max-idle=30mAdvice about initialisation actions:
- The script runs on every node, including the ones autoscaling adds. It must be idempotent and fast: a five-minute script multiplies each new node's start-up time by five.
- Use
get_metadata_value attributes/dataproc-roleto tell master from worker. set -euxo pipefailmakes the script fail loudly instead of leaving the node half-configured.- If there are many dependencies, it is better to build a custom image than to install them on every start-up.
About the image versions: each Dataproc version packages a specific set of Spark, Hadoop, Python and the operating system. The 2.2-debian12 series, for example, brings Spark 3.5 and Python 3.11.
| Practice | Consequence |
|---|---|
| Not stating a version | Google picks the most recent one: your job can break on its own |
Stating the series (2.2-debian12) |
Automatic minor updates within the series. A reasonable balance |
Stating the exact version (2.2.28-debian12) |
Full reproducibility. Recommended in critical production |
For AlpinaShop: a pinned series in development, an exact version in the scheduled workflows. And in the workflow file, versioned in Git.
- Dataproc Serverless for Spark
Although an ephemeral cluster is far better than a permanent one, you still have to size it: how many workers, which machine, how much memory per executor. Dataproc Serverless removes that decision: you submit the Spark job and Google takes care of everything.
gcloud dataproc batches submit pyspark \
gs://alpinashop-datalake/jobs/cesta_media.py \
--batch=cesta-media-$(date +%Y%m%d-%H%M%S) \
--region=europe-west1 \
--version=2.2 \
--service-account="$SA" \
--subnet=sn-datos-euw1 \
--deps-bucket=gs://alpinashop-datalake \
--properties="spark.executor.instances=4,\
spark.dynamicAllocation.enabled=true,\
spark.dynamicAllocation.maxExecutors=20" \
--labels=entorno=produccion,equipo=datos,centro-coste=analitica \
-- --date-from=2025-01-01 --date-to=2026-03-31There is no clusters create. There is no --max-idle because there is nothing to switch off. There is no --num-workers.
| Aspect | Dataproc with a cluster | Dataproc Serverless |
|---|---|---|
| Management | You create and destroy clusters | None |
| Start-up | 90-120 s | 30-60 s |
| Sizing | You choose it | Automatic |
| Billing | Per VM per hour | Per DCU (data compute units) while it lasts |
| Components | The whole ecosystem: Hive, HBase, Presto, Jupyter | Spark only |
| Initialisation actions | Yes | No; custom container images are used instead |
| Spot VMs | Yes, very cheap | Not applicable |
| Interactive session | Notebook on the cluster | Serverless interactive sessions |
The criterion: if your job is pure Spark and you need neither Hive nor HBase nor a long-lived cluster, start with Serverless. It is less to administer and less to leave switched on by accident. Use a cluster when you need ecosystem components, long interactive sessions, or when the Spot VM discount on very large workloads outweighs the management.
For AlpinaShop, the reasoned decision is: the basket analysis goes to Serverless, because it is pure Spark, it runs monthly and nobody wants to remember to switch anything off. The alpinashop-spark cluster is kept solely as an exploration environment with Jupyter, with --max-idle=30m, and it is destroyed when it goes unused for a month.
- Dataproc, Serverless and Dataflow: the decision table
| Criterion | Dataproc (cluster) | Dataproc Serverless | Dataflow |
|---|---|---|---|
| Programming model | Spark / Hadoop / Hive | Spark | Apache Beam |
| Batch | Yes | Yes | Yes |
| Streaming | Structured Streaming | Limited | Its strong point |
| Start-up | ~2 min | ~40 s | ~2 min (batch) |
| Infrastructure | You define it | None | None |
| Autoscaling | With a policy | Automatic | Automatic |
| Cost at rest | The cluster, if you leave it | Zero | Zero (unless streaming is active) |
| Exactly-once semantics | You have to build it | Same | As standard |
| Windows and event time | Manual | Manual | A complete model |
| Library ecosystem | Enormous (MLlib, pandas…) | Large | Limited to Beam |
| Portability | High (Spark runs everywhere) | Medium | High (Beam has several runners) |
| Choose it if… | You have Spark code, you need Hive/HBase, the team knows Spark | Pure Spark and you do not want to manage anything | Streaming, event time, new pipelines |
AlpinaShop's final decision comes out like this:
- Streaming of orders and visits → Dataflow. Beam's model of windows and watermarks has no comfortable equivalent in Spark, and the Pub/Sub to BigQuery templates solve the base case with no code.
- New batch loads and transformations → Dataflow, for consistency with the above and because the team already has it set up.
- Algorithmic analysis: basket, co-occurrence, future models with MLlib → Dataproc Serverless.
- Interactive exploration over the lake → a notebook on
alpinashop-spark, with--max-idle. - Transformation expressible in SQL over data already in BigQuery → BigQuery, without moving anything.
- Notebooks and Spark SQL over BigQuery
With --optional-components=JUPYTER and --enable-component-gateway, the Dataproc console shows a link to JupyterLab protected by IAM. No opening ports and no SSH tunnels: whoever has the right role gets in, and whoever does not, does not.
The BigQuery connector for Spark lets you read and write tables directly:
from pyspark.sql import SparkSession, functions as F
spark = (SparkSession.builder
.appName("exploracion-lucia")
.config("temporaryGcsBucket", "alpinashop-datalake")
.getOrCreate())
# Read a whole table: the connector uses the Storage Read API,
# which reads in parallel and in columnar format. It does not export to files.
products = (spark.read.format("bigquery")
.option("table", "alpinashop-datos.alpinashop_analitica.productos")
.load())
# BETTER: delegate the filter to BigQuery so that less data comes across
lines = (spark.read.format("bigquery")
.option("table", "alpinashop-datos.alpinashop_analitica.lineas_pedido")
.option("filter", "fecha_pedido >= '2026-01-01'") # runs in BigQuery
.load())
# From here on, ordinary Spark SQL
lines.createOrReplaceTempView("lineas")
products.createOrReplaceTempView("productos")
summary = spark.sql("""
SELECT p.categoria,
COUNT(DISTINCT l.pedido_id) AS pedidos,
ROUND(SUM(l.importe_linea), 2) AS ventas_eur
FROM lineas l
JOIN productos p ON p.sku = l.sku
GROUP BY p.categoria
ORDER BY ventas_eur DESC
""")
summary.show(truncate=False)The filter option is what makes the difference: it runs in BigQuery, not in Spark. Without it, the connector would bring the whole table over the network for Spark to filter, paying for the full read in BigQuery and wasting time. It is the same principle we saw with EXTERNAL_QUERY in 04-01: filter as close to the source as possible.
And the obvious question: if you can do this in Spark SQL, why not do it in BigQuery directly? Almost always you should. The connector makes sense when the SQL result feeds an MLlib algorithm, when you cross BigQuery tables with lake files that have not been loaded, or when the Spark code already exists. For a pure aggregation, BigQuery is faster and cheaper.
- Migrating an on-premises Hadoop to Google Cloud
This is the scenario for which Dataproc is most justified. Suppose AlpinaShop absorbs a competitor with a 30-node Hadoop cluster.
What is preserved:
- The Spark, Hive and PySpark code: the vast majority works unchanged.
- The Hive queries: Dataproc includes Hive, and the metastore can be migrated.
- The data formats: Parquet, ORC and Avro are identical.
- The Oozie workflows, although it is worth replacing them with Composer (04-06).
What must be rethought, without exception:
| On-premises element | On Google Cloud | Why it changes |
|---|---|---|
| Permanent HDFS | Cloud Storage | It is the fundamental change: without it there are no ephemeral clusters and no saving |
| One giant shared cluster | Several small clusters, one per workload | Each team with its own version and its own budget; no queues and no noisy neighbours |
| Sized for the peak | Autoscaling + Spot | You were paying for the peak 24 hours a day |
| Kerberos | IAM and service accounts | The cloud identity model (03-04) |
| Local Hive metastore | Managed Dataproc Metastore | It outlives the ephemeral clusters: it is the piece that makes them possible |
| Oozie / cron | Cloud Composer or Workflows | 04-06 |
| Impala / Presto for queries | BigQuery | Usually the biggest leap in performance and simplicity |
| Flume / Kafka for ingestion | Pub/Sub (04-04) or Managed Kafka | Managed |
The recommended strategy, and the only one that usually works out, is phased:
- Copy the data to Cloud Storage with the Storage Transfer Service or
hadoop distcp, touching nothing else. The on-premises cluster carries on working. - Stand up Dataproc Metastore and register the tables pointing at
gs://instead ofhdfs://. - Run the existing jobs on Dataproc against the data already in Cloud Storage, comparing results with those from the old cluster. This dual-run phase is non-negotiable: it is the only way to prove that the numbers match.
- Switch off the on-premises cluster when the comparison holds for several weeks.
- Only then, optimise: move the Hive queries to BigQuery, the streaming jobs to Dataflow, adopt Serverless.
The classic mistake is attempting step 5 at the same time as step 3, that is, migrating and modernising simultaneously. When the numbers do not match, nobody knows whether it is because of the migration or the rewrite, and the project stalls for months.
- Cost and Spot VMs
Dataproc bills two things:
- The underlying VMs (Compute Engine, disks and network), at the normal rate.
- A Dataproc management fee, of the order of $0.01 per vCPU per hour, verifiable in the official documentation.
In other words, Dataproc's overhead compared with building Hadoop yourself on VMs is small: you pay very little for not administering anything.
An example with alpinashop-spark (1 master + 2 workers + 2 secondaries, all n2-standard-4, 20 vCPU in total), as an order of magnitude:
| Scenario | Hours per month | Approximate cost |
|---|---|---|
| Permanent 24×7 cluster | 720 | Of the order of €1,400 |
Cluster with --max-idle=30m, 2 h of use a day |
~75 | Of the order of €150 |
| Ephemeral cluster per workflow, 20 min a day | ~10 | Of the order of €20 |
| With Spot secondaries instead of standard | ~10 | Of the order of €12 |
| Dataproc Serverless, the same monthly job | ~1 | Cents |
Verify current prices in the official documentation; what matters here is the factor of 100 between the first row and the third. That factor is the whole lesson.
Saving levers, in order of impact:
--max-idle, always. A cluster forgotten over a four-day break costs more than a year of Serverless.- Ephemeral clusters per workflow. It eliminates the problem at the root.
- Spot secondary workers. Discounts of up to 80 %, with the caveat that they do not contribute HDFS and can disappear.
- Serverless for occasional jobs.
- Right-sized disks. With the data in Cloud Storage, HDFS is only shuffle space: 200 GB per worker is more than enough for almost everything.
- A consistent region. Cluster and buckets in
europe-west1: reading data from another region costs egress and latency. - Committed use discounts only if you end up with a permanent cluster, which this section suggests avoiding.
- Columnar formats. Parquet with files of 128-512 MB reads far less and avoids the small-files problem, which is the biggest performance killer in Spark over object storage.
Common Mistakes and Tips
Creating a permanent cluster out of habit. It is the mental inheritance of on-premises Hadoop and it is the most expensive mistake. If the cluster has no --max-idle, it should not exist.
Writing important data to HDFS. It disappears when the cluster is destroyed. HDFS in Dataproc is working memory, not storage.
Not pinning the image version. Google updates the default version and a job that used to work stops working without anybody having touched the code.
Using collect() on a large DataFrame. It brings all the data back to the driver, which is a single machine. It is the number one cause of OutOfMemoryError in Spark. Use show(), take(n) or write to a file.
Forgetting cache() when a DataFrame is used several times. Spark recomputes the whole chain on every action. And the opposite mistake: caching everything that moves, until memory fills up and it starts spilling to disk.
Leaving spark.sql.shuffle.partitions at 200. With small data it generates 200 tiny tasks with more overhead than work; with large data, enormous partitions that do not fit in memory.
Lots of small files in the lake. Ten thousand 1 MB files are far slower to read than twenty of 500 MB, because every open has latency. Compact them.
Putting the data in one region and the cluster in another. Billed cross-region egress and added latency on every read.
Tip: use Serverless by default. Start there and create a cluster only when you discover you need something Serverless does not give you. It is the path with the least operational debt.
Tip: --dry-run does not exist, but the subset does. Before launching a job over two years, run it over one week. Logic errors show up just the same and cost a hundred times less.
Tip: always look at the Spark interface. The component gateway gives access to the Spark UI, where you can see the stages, the tasks and — most useful of all — the distribution of times across tasks. If one task takes a hundred times longer than the median, you have data skew, exactly as in Dataflow.
Exercises
Exercise 1: an ephemeral cluster with self-destruction
Create a cluster called alpinashop-spark-pruebas in europe-west1 that: lives in sn-datos-euw1 with no public IP, uses the sa-dataproc-analitica service account, has 1 n2-standard-2 master and 2 n2-standard-2 workers, adds 2 Spot secondary workers, pins the 2.2-debian12 image, self-destructs after 15 minutes of inactivity and in any case after 2 hours of life, and carries AlpinaShop's standard labels. Then check its status and delete it explicitly.
Exercise 2: a PySpark returns job
Write a PySpark job that reads gs://alpinashop-datalake/pedidos/cabeceras/ and gs://alpinashop-datalake/pedidos/lineas/ in Parquet, and computes per product category: the number of orders with at least one return, the return rate over the total orders in that category, and the average amount returned. The result must be written to alpinashop-datos.alpinashop_analitica.devoluciones_categoria. Apply at least two of the optimisations seen in the lesson and explain why you apply them.
Exercise 3: choosing the tool
For each of these five AlpinaShop assignments, choose between BigQuery, Dataflow, Dataproc with a cluster, Dataproc Serverless or bq load, and justify the choice in two or three sentences:
- Loading a 4 GB Parquet file from the ERP into a BigQuery table every night, with no transformation whatsoever.
- Computing the monthly sales ranking by category from tables that are already in
alpinashop_analitica. - Processing the
pedidos-nuevosevents as they arrive, grouping them by the event's hour and tolerating 90 minutes of delay. - Training a recommendation model weekly with MLlib's ALS over two years of purchase history.
- Running the PySpark code inherited from the absorbed competitor, 15,000 lines that use Hive and user-defined functions, while it is decided what to do with it.
Solutions
Solution 1
SA="[email protected]"
gcloud dataproc clusters create alpinashop-spark-pruebas \
--region=europe-west1 \
--zone=europe-west1-b \
--subnet=sn-datos-euw1 \
--no-address \
--service-account="$SA" \
--master-machine-type=n2-standard-2 \
--master-boot-disk-size=100GB \
--num-workers=2 \
--worker-machine-type=n2-standard-2 \
--worker-boot-disk-size=100GB \
--num-secondary-workers=2 \
--secondary-worker-type=spot \
--image-version=2.2-debian12 \
--max-idle=15m \
--max-age=2h \
--labels=entorno=desarrollo,equipo=datos,centro-coste=analitica,aplicacion=analitica# Check the status
gcloud dataproc clusters describe alpinashop-spark-pruebas --region=europe-west1 \
--format="yaml(status.state, config.lifecycleConfig, config.gceClusterConfig.internalIpOnly)"
gcloud dataproc clusters list --region=europe-west1 \
--format="table(clusterName, status.state, config.softwareConfig.imageVersion)"
# Explicit deletion, without waiting for max-idle
gcloud dataproc clusters delete alpinashop-spark-pruebas --region=europe-west1 --quietThe two lifecycle options are complementary and it is worth setting both:
--max-idle=15m: it is destroyed if nobody uses it for 15 minutes. It covers the "I went to lunch" case.--max-age=2h: it is destroyed after 2 hours no matter what. It covers the "a job in an infinite loop keeps the cluster busy and--max-idlenever fires" case. That is the failure that produces weekend invoices.
--no-address with --subnet=sn-datos-euw1 requires Private Google Access to be enabled on that subnet (03-01) or the nodes will not be able to download packages or talk to the APIs. If the cluster stays in CREATING and then fails, that is the first cause to check.
Solution 2
"""devoluciones.py -- Return rate by category."""
from pyspark.sql import SparkSession, functions as F
PROJECT, DATASET = "alpinashop-datos", "alpinashop_analitica"
LAKE = "gs://alpinashop-datalake"
spark = (SparkSession.builder
.appName("alpinashop-devoluciones")
.config("temporaryGcsBucket", "alpinashop-datalake")
.getOrCreate())
# OPTIMISATION 1: select only the necessary columns when reading.
# Parquet is columnar: the columns not requested are never read from the bucket.
headers = (spark.read.parquet(f"{LAKE}/pedidos/cabeceras/")
.select("pedido_id", "estado", "fecha_pedido", "total_pedido"))
lines = (spark.read.parquet(f"{LAKE}/pedidos/lineas/")
.select("pedido_id", "sku", "importe_linea"))
# The catalogue is read from BigQuery, not from the lake
products = (spark.read.format("bigquery")
.option("table", f"{PROJECT}.{DATASET}.productos")
.option("filter", "activo = true")
.load()
.select("sku", "categoria"))
# OPTIMISATION 2: broadcast the catalogue (thousands of rows, it fits in memory).
# It completely avoids the JOIN shuffle with the lines table, which is
# the big one. Without this, both sides would be repartitioned by 'sku'.
lines_cat = lines.join(F.broadcast(products), on="sku", how="inner")
# A returned order is returned in full: we flag it at header level
hdr = headers.withColumn(
"is_returned", F.when(F.col("estado") == "devuelto", 1).otherwise(0)
)
# OPTIMISATION 3: cache, because the joined DataFrame is used twice
detail = lines_cat.join(F.broadcast(hdr), on="pedido_id", how="inner")
detail.cache()
# Distinct orders per category (one order can touch several categories)
per_category = (
detail.groupBy("categoria")
.agg(
F.countDistinct("pedido_id").alias("pedidos_totales"),
F.countDistinct(
F.when(F.col("is_returned") == 1, F.col("pedido_id"))
).alias("pedidos_devueltos"),
F.round(
F.avg(F.when(F.col("is_returned") == 1, F.col("importe_linea"))), 2
).alias("importe_medio_devuelto"),
)
.withColumn(
"tasa_devolucion_pct",
F.round(100 * F.col("pedidos_devueltos") / F.col("pedidos_totales"), 2),
)
.orderBy(F.desc("tasa_devolucion_pct"))
)
(per_category.write.format("bigquery")
.option("table", f"{PROJECT}.{DATASET}.devoluciones_categoria")
.option("writeMethod", "direct")
.mode("overwrite")
.save())
per_category.show(truncate=False)
detail.unpersist()
spark.stop()The optimisations and their justification:
- Early column projection. Parquet is columnar; asking for four columns instead of twenty proportionally reduces the bytes read from the bucket and the executors' memory. It is the exact equivalent of not doing
SELECT *in BigQuery. broadcast()on bothJOINs. The product catalogue is a few thousand rows and the period's headers are also small compared with the lines. Broadcasting them avoids shuffling the big table over the network, which is the dominant cost.cache()ondetail. Although in this final version there is only one aggregation, as soon as a second calculation is added (by country, by month) Spark would recompute the whole chain. It is the correct preparation; withunpersist()at the end so as not to hold memory.- Filter delegated to BigQuery (
option("filter", "activo = true")). It runs over there and fewer rows come across.
A methodological note that has to be made explicit when presenting the result: an order with products from three categories counts as an order in all three, so the per-category figures do not add up to the total number of orders. That is correct for measuring a rate per category, but it has to be stated in the report or someone will subtract and the numbers will not add up for them.
Solution 3
1. Loading a 4 GB Parquet file every night with no transformation → bq load.
There is no transformation, so no processing engine is needed. Batch loading into BigQuery is free, Parquet carries the schema built in and there is nothing to size. Using Dataflow or Spark here would mean paying for compute to make a copy. A single command, orchestrated in 04-06.
2. Monthly ranking over tables already in alpinashop_analitica → BigQuery.
The data is already there and the operation is an aggregation with a window function, exactly what we did in 04-01 with RANK() OVER and QUALIFY. Taking the data out of BigQuery to process it elsewhere and putting it back in is the classic anti-pattern: read cost, compute cost, write cost and latency, for a worse result. If it has to be refreshed often, a materialized view.
3. pedidos-nuevos events by the event's hour with 90 minutes of tolerance → Dataflow.
This is literally the use case Beam's model exists for: unbounded streaming, grouping by event time, fixed windows, watermark and allowed_lateness. Spark's Structured Streaming could do it, but with a less expressive time model and without native Pub/Sub integration. On top of that, Dataflow gives exactly-once semantics towards BigQuery with no extra work.
4. Training MLlib's ALS weekly → Dataproc Serverless. ALS is a distributed iterative algorithm implemented in MLlib, with no equivalent in Beam or in pure SQL. It is Spark from start to finish. Serverless rather than a cluster because it runs once a week: nobody wants to maintain or remember to switch off a cluster that works for an hour every seven days. (The next step, serving that model in production, is Vertex AI territory in module 5.)
5. 15,000 inherited lines of PySpark with Hive and UDFs → Dataproc with a cluster. Here the practical constraint rules: the code exists, it works and it uses Hive, which Serverless does not include. The priority is for it to carry on working with the minimum of change, so a Dataproc cluster with Dataproc Metastore, running the code as-is against the data already copied to Cloud Storage. Modernisation — moving queries to BigQuery, streaming to Dataflow — is a later and separate phase, never simultaneous with the migration: if the numbers do not match, you have to be able to tell whether it is the move or the rewrite.
Conclusion
You have toured the world of Hadoop and Spark with exactly the right perspective: what they are, what problem they solved — processing more data than fits on one machine, with machines that fail — why Spark displaced MapReduce by keeping the intermediate data in memory, and why in 2026 they still matter even though BigQuery and Dataflow exist: there is code already written, there are people who know how to use it, there are libraries such as MLlib with no equivalent, and there is an open ecosystem that provides real portability.
You have seen what Dataproc adds: clusters in under two minutes, configured and coherent, integrated with Cloud Storage, BigQuery, IAM and Logging. You know the anatomy — master, primary workers that hold up HDFS, secondary workers that only contribute compute and can therefore be Spot — and you have created alpinashop-spark inside sn-datos-euw1, with no public IP, with the sa-dataproc-analitica account, with the image pinned, with Jupyter reachable through the component gateway and with an autoscaling policy that only scales the secondaries and removes them gracefully.
But the important thing in this lesson is not a command, it is a conceptual inversion: the ephemeral cluster over storage in Cloud Storage. Because the internal network makes reading from gs:// comparable to reading from a local disk, the data no longer needs to live on the cluster. And if the data does not live on the cluster, the cluster can die. From that come --max-idle, --max-age, workflows with a managed cluster that is born and dies with the job, and the hundredfold difference on the invoice between a permanent cluster and an ephemeral one.
You have written the PySpark job that answers Lucía's question: the average basket per month and country with its median next to it, and the matrix of products bought together with support, confidence and lift — the metric that corrects for popularity and stops you concluding that everybody buys socks with everything — applying broadcast to avoid shuffle, cache to avoid recomputation, a cap on lines per order so that rare corporate orders do not dominate the calculation, and the sku_a < sku_b trick to count each pair only once. You know the initialisation actions, their risks and why pinning the image version is not a quirk.
You have met Dataproc Serverless, which removes even the sizing decision, and you have set AlpinaShop's criterion: the basket analysis on Serverless, the cluster only as an exploration environment with Jupyter and self-destruction. And you have the complete decision table between Dataproc, Serverless and Dataflow, with streaming and event time on Beam's side, algorithms and ecosystem on Spark's side, and SQL over already-loaded data on BigQuery's side, without moving anything. You know how an on-premises Hadoop is migrated in phases, with the golden rule of not modernising and migrating at the same time. And you know where the money is: --max-idle, ephemeral clusters, Spot, Serverless, right-sized disks and large files in a columnar format.
There is a promise that has been outstanding for two lessons. Everything you have built — the Dataflow pipeline, the Spark job, the BigQuery tables — works over data that is already somewhere. But the shop is still an island: when a customer confirms an order, the Flask application has to tell the warehouse to prepare it, billing to issue the invoice, the mail service to send the confirmation, and now analytics as well. If it does that by calling all four one after the other, the sale hangs waiting for the slowest, and if one fails, it is not clear what has happened to the other three. It is a fragile design that breaks on exactly the busiest sales day of the year.
In 04-04, Cloud Pub/Sub, we will break that coupling. We will finally create the pedidos-nuevos topic and the sub-almacen, sub-facturacion and sub-analitica subscriptions; you will understand messaging's real guarantees — at-least-once delivery, no guaranteed order — and why that forces your consumers to be idempotent; you will see dead letter topics, retries with exponential backoff, message replay with seek, filters by attribute and the direct BigQuery subscriptions that ingest without writing a line of code. And we will finally connect, properly, the alpinashop-catalogo bucket notifications that were promised back in 02-02.
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
