alpinashop_analitica already exists and the order history is sitting inside it. But that history got there in a way that cannot be repeated every night: Lucía ran a federated query against the PostgreSQL replica by hand. It worked once. The question is what happens tomorrow, and the day after, and the day someone asks to see the campaign's sales today rather than tomorrow.
The naive answer is to write a script. A sincronizar.py file that connects to Cloud SQL, reads the day's orders, transforms them and inserts them into BigQuery, launched by a cron on a virtual machine. That script will work for weeks, and then it will fail. It will fail because the data grows and the script takes five hours; because the VM restarts halfway through and nobody knows whether it wrote half the rows; because one order with an odd character throws an exception and the rest of the dump is lost; because nobody knows whether it ran yesterday; and because the day someone asks for real-time data, the whole design has to be thrown away.
Cloud Dataflow is Google's managed service for running data pipelines written with Apache Beam. Its concrete promise is twofold: you write the transformation logic once and Google takes care of parallelism, retries, scaling and consistency; and that same code serves both to process two years of history and to process orders as they come in.
In this lesson you are going to understand Beam's model, write a batch pipeline that loads AlpinaShop's history, test it locally, launch it on Dataflow, and then get into the part that genuinely separates serious data processing from the amateur kind: time. What "the 10:00 sales" means when an event generated at 10:00 arrives at 10:07, and how you answer that without lying.
Contents
- What a data pipeline is and why a script is not enough
- Apache Beam: the unified model
- The four concepts:
Pipeline,PCollection,PTransform, runner - The transformations you will use 90 % of the time
- The first batch pipeline, explained line by line
- Running locally with
DirectRunner - Managed execution with
DataflowRunner - Time in streaming: event versus processing
- Watermarks, windows, triggers and late data
- The
pedidos-nuevosstreaming pipeline - Dataflow templates: the practical route
- Autoscaling and Dataflow Prime
- Monitoring, parallelism and data skew
- Cost: what exactly you pay for
- When Dataflow is not the answer
- What a data pipeline is and why a script is not enough
A data pipeline is a declared sequence of operations that carry data from a source to a destination, transforming it along the way. The important adjective is declared: you describe what you want to happen, not how the work is shared out across machines.
Let us compare the script on a VM with the managed pipeline honestly:
| Aspect | Script on a VM | Pipeline on Dataflow |
|---|---|---|
| Parallelism | Whatever you program yourself, by hand, with threads | Automatic: shared out across workers |
| Scaling | Change the VM and restart | Horizontal and automatic during execution |
| Machine failure | The whole process is lost | The affected fragment is retried |
| Record failure | An exception that brings the process down | Diverted to an error output and it carries on |
| Write semantics | Whatever you manage to achieve | Exactly once on the native connectors |
| State after a restart | Unknown | Managed by the service |
| Batch and streaming | Two different programs | The same code |
| Cost at rest | The VM switched on permanently | Zero: there is nothing switched on |
| Observability | The print calls you happened to add |
Graph, metrics and logs built in |
The row that hurts most in practice is record failure. A dump of 800,000 orders in which order 412,337 has an amount with a comma instead of a full stop must not lose the remaining 387,663 orders. A badly built script loses them; a well-built one demands considerable error-handling effort that in Beam comes as standard with tagged outputs.
And the cost-at-rest row has its nuance: Dataflow in batch mode costs zero when it is not running, but a streaming pipeline is permanently switched on and bills without pause. We will come back to that in section 14, because it is the most common surprise.
- Apache Beam: the unified model
Apache Beam is an open source programming model — donated by Google to the Apache Software Foundation — for defining data pipelines that are then run on different engines.
The key word is unified. Before Beam, batch processing and streaming processing were separate worlds with separate tools, and companies maintained two implementations of the same business logic: one for the history and one for real time, which inevitably diverged and produced different numbers. Beam starts from the idea that a batch is simply a bounded flow and a stream is an unbounded flow, and that the transformation logic is the same in both cases.
flowchart LR
subgraph SDK["Apache Beam SDK"]
P["Pipeline code<br/>Python / Java / Go"]
end
subgraph Runners["Runners"]
D["DirectRunner<br/>local, testing"]
DF["DataflowRunner<br/>Google Cloud"]
FL["FlinkRunner"]
SP["SparkRunner"]
end
P --> D
P --> DF
P --> FL
P --> SP
The runner is the engine that executes the pipeline. The same Python file can run on your laptop with DirectRunner, on Dataflow with DataflowRunner, or on a Flink or Spark cluster. That reduces lock-in with the provider: if AlpinaShop ever had to leave Google Cloud, its pipeline logic would travel.
Beam has SDKs for Java, Python and Go. We will use Python, consistent with the rest of AlpinaShop's application, which is already Flask.
# Installing the SDK with the Google Cloud dependencies
python -m venv venv-beam
source venv-beam/bin/activate
pip install 'apache-beam[gcp]==2.64.0'Pinning the version is deliberate: the Dataflow workers will use exactly the SDK version you launch the pipeline with, and differences between versions are a classic source of failures that only show up in the cloud.
- The four concepts:
Pipeline, PCollection, PTransform, runner
Pipeline, PCollection, PTransform, runnerPipeline is the object that holds the complete graph. It is built, declared and executed. Nothing happens while you write it: you are drawing a blueprint.
PCollection is a distributed and immutable dataset. It is not a Python list: it can have zero elements or trillions, it can be spread across a hundred machines, and it cannot be modified. Every transformation produces a new PCollection. Immutability is what makes it possible to retry a failed fragment without corrupting anything.
A PCollection can be:
- Bounded: it has a known end. A file, a table. This is the batch case.
- Unbounded: it never ends. A Pub/Sub topic. This is the streaming case.
PTransform is an operation that takes one or more PCollection and produces one or more PCollection. It is applied with the | operator, which in Beam is overloaded to mean "apply this transformation".
Runner is the execution engine, already covered.
The syntax, which looks odd the first time:
The >> operator associates a name with the step. It is not decorative: that name is the one that appears in the graph in the Dataflow console and in the metrics. A pipeline with steps called Map(<lambda at main.py:34>) is impossible to debug in production. Name every step, always.
flowchart TD
A["PCollection: raw text lines<br/>gs://alpinashop-catalogo/exportaciones/..."]
B["PTransform: ParseCSV<br/>ParDo"]
C["PCollection: order dictionaries"]
D["PTransform: ValidateAndClean<br/>ParDo with multiple outputs"]
E["PCollection: valid orders"]
F["PCollection: rejected orders"]
G["PTransform: WriteToBigQuery"]
H["PTransform: WriteToText<br/>quarantine in the bucket"]
A --> B --> C --> D
D --> E --> G
D --> F --> H
That graph is exactly the one you are going to write in section 5. Notice that the error branch is part of the design, not an afterthought.
- The transformations you will use 90 % of the time
| Transformation | What it does | Example at AlpinaShop |
|---|---|---|
beam.Map(f) |
Applies f to each element; returns one |
Turn a CSV line into a dictionary |
beam.FlatMap(f) |
Applies f; returns 0, 1 or N elements |
Break an order down into its lines |
beam.Filter(f) |
Keeps those that satisfy f |
Discard cancelled orders |
beam.ParDo(DoFn) |
The general form: a class with state, lifecycle and multiple outputs | Validate and separate valid from rejected |
beam.GroupByKey() |
Groups (key, value) pairs by key |
Group lines by sku |
beam.CombinePerKey(f) |
Groups and reduces with an associative function | Sum sales by sku |
beam.CoGroupByKey() |
Joins several PCollection by key (Beam's JOIN) |
Cross orders with products |
beam.Keys() / beam.Values() |
Extracts keys or values | |
beam.Distinct() |
Removes duplicates | Unique sessions |
beam.io.ReadFromText / WriteToText |
Reading/writing files, including Cloud Storage | |
beam.io.ReadFromPubSub |
Reads a topic or subscription | The pedidos-nuevos pipeline |
beam.io.WriteToBigQuery |
Writes to a table | The destination of everything |
Two clarifications that prevent important conceptual mistakes:
GroupByKey versus CombinePerKey. GroupByKey brings all the values for a key to a single machine and materialises them in memory. If one sku has three million lines, that machine can blow up. CombinePerKey with an associative and commutative function (sum, max, min) performs partial aggregation on each worker before moving anything over the network: each worker sums its own share and only the partial results travel. Whenever you can use CombinePerKey, use it; it is orders of magnitude more efficient and it does not break with hot keys.
Map versus ParDo. Map is syntactic sugar over ParDo. Use Map for simple, stateless transformations; use ParDo with a DoFn class when you need expensive initialisation (opening an API client once per worker, not once per element), your own metrics, or several outputs.
- The first batch pipeline, explained line by line
The concrete goal: read the order export CSV files sitting in gs://alpinashop-catalogo/exportaciones/2026/03/14/, validate them, clean them, compute a sales aggregate per SKU and write two things into alpinashop_analitica: the clean lines into lineas_pedido and the aggregate into a new table ventas_diarias_sku. Faulty records go to a quarantine file in the bucket.
"""
AlpinaShop batch pipeline: CSV exports -> BigQuery.
Local run:
python pipeline_pedidos.py --date 2026-03-14
Run on Dataflow:
python pipeline_pedidos.py --date 2026-03-14 --runner DataflowRunner ...
"""
import argparse
import csv
import io
import logging
from datetime import datetime
from decimal import Decimal, InvalidOperation
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions
PROJECT = "alpinashop-datos"
DATASET = "alpinashop_analitica"
BUCKET = "alpinashop-catalogo"Everything above is ordinary Python. The uppercase constants go at the top so that changing project does not force you to hunt for strings through the file.
class ParseCSVLine(beam.DoFn):
"""Turns a CSV text line into a dictionary.
It is implemented as a DoFn and not as a Map because we need
our own counters and because we want two outputs: valid and rejected.
"""
REJECTED_OUTPUT = "rejected"
def __init__(self):
# Beam counters are aggregated across all workers
# and are visible in the Dataflow console. They are the correct way
# to instrument a pipeline; print statements are useless.
self.counter_ok = beam.metrics.Metrics.counter("parsing", "rows_ok")
self.counter_ko = beam.metrics.Metrics.counter("parsing", "rows_rejected")
def process(self, line):
try:
fields = next(csv.reader(io.StringIO(line), delimiter=","))
except Exception as exc:
self.counter_ko.inc()
yield beam.pvalue.TaggedOutput(
self.REJECTED_OUTPUT,
{"line": line, "reason": f"unreadable_csv: {exc}"},
)
return
if len(fields) != 8:
self.counter_ko.inc()
yield beam.pvalue.TaggedOutput(
self.REJECTED_OUTPUT,
{"line": line, "reason": f"expected 8 columns, found {len(fields)}"},
)
return
yield {
"pedido_id": fields[0].strip(),
"linea_num": fields[1].strip(),
"fecha_pedido": fields[2].strip(),
"sku": fields[3].strip().upper(),
"cantidad": fields[4].strip(),
"precio_unitario": fields[5].strip().replace(",", "."),
"descuento_linea": fields[6].strip().replace(",", "."),
"estado": fields[7].strip().lower(),
}Important points about this class:
yieldinstead ofreturn. ADoFnis a generator: it can emit zero, one or many elements per input. Areturnwith a value does not work the way you expect.TaggedOutputmarks the element so that it leaves through a different branch of the graph. It is Beam's mechanism for "this went wrong but the pipeline carries on". Without it, an exception retries the bundle four times and then kills the whole job.replace(",", ".")normalises the decimals from the Spanish ERP, which exports89,90. It is exactly the kind of real-world dirt that lives in any export.- The counters (
Metrics.counter) are pushed to the Dataflow console and let you answer "how many rows were rejected last night?" without opening a log.
class ValidateOrder(beam.DoFn):
"""Converts types and applies business rules."""
REJECTED_OUTPUT = "rejected"
def process(self, row):
reasons = []
# 1) Date
try:
order_date = datetime.strptime(row["fecha_pedido"], "%Y-%m-%d").date()
except ValueError:
reasons.append("invalid date")
order_date = None
# 2) Quantity: strictly positive integer
try:
quantity = int(row["cantidad"])
if quantity <= 0:
reasons.append("non-positive quantity")
except ValueError:
reasons.append("non-numeric quantity")
quantity = None
# 3) Amounts as Decimal, NEVER as float (see 04-01)
try:
price = Decimal(row["precio_unitario"])
discount = Decimal(row["descuento_linea"] or "0")
if price < 0 or discount < 0:
reasons.append("negative amount")
except InvalidOperation:
reasons.append("non-numeric amount")
price = discount = None
# 4) SKU in the format of the AlpinaShop catalogue
if not row["sku"] or len(row["sku"]) < 4:
reasons.append("sku missing or too short")
if reasons:
yield beam.pvalue.TaggedOutput(
self.REJECTED_OUTPUT,
{"line": str(row), "reason": "; ".join(reasons)},
)
return
amount = (price * quantity) - discount
yield {
"pedido_id": row["pedido_id"],
"linea_num": int(row["linea_num"]),
"fecha_pedido": order_date.isoformat(),
"sku": row["sku"],
"cantidad": quantity,
"precio_unitario": str(price), # BigQuery accepts NUMERIC as a string
"descuento_linea": str(discount),
"importe_linea": str(amount),
}A detail that is expensive to discover in production: NUMERIC values are passed to WriteToBigQuery as strings, not as float. If you convert to float to serialise, you have reintroduced the floating-point error we avoided so carefully in 04-01.
Now the complete pipeline:
def build_pipeline(pipeline, date):
input_path = f"gs://{BUCKET}/exportaciones/{date.replace('-', '/')}/pedidos-*.csv"
quarantine_path = f"gs://{BUCKET}/cuarentena/{date}/rechazos"
# 1) READ: each CSV line is one element of the PCollection
raw = (
pipeline
| "ReadCSV" >> beam.io.ReadFromText(input_path, skip_header_lines=1)
)
# 2) PARSE with two outputs
parsed = (
raw
| "ParseCSV" >> beam.ParDo(ParseCSVLine()).with_outputs(
ParseCSVLine.REJECTED_OUTPUT, main="valid"
)
)
# 3) VALIDATE, also with two outputs
validated = (
parsed.valid
| "ValidateOrder" >> beam.ParDo(ValidateOrder()).with_outputs(
ValidateOrder.REJECTED_OUTPUT, main="clean"
)
)
clean_lines = validated.clean
# 4) WRITE the detail to BigQuery
(
clean_lines
| "WriteLines" >> beam.io.WriteToBigQuery(
table=f"{PROJECT}:{DATASET}.lineas_pedido",
schema="pedido_id:STRING,linea_num:INTEGER,fecha_pedido:DATE,"
"sku:STRING,cantidad:INTEGER,precio_unitario:NUMERIC,"
"descuento_linea:NUMERIC,importe_linea:NUMERIC",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
additional_bq_parameters={
"timePartitioning": {"type": "DAY", "field": "fecha_pedido"},
"clustering": {"fields": ["sku"]},
},
)
)
# 5) AGGREGATE sales per SKU with CombinePerKey (not GroupByKey)
(
clean_lines
| "KeyBySKU" >> beam.Map(
lambda f: ((f["fecha_pedido"], f["sku"]), Decimal(f["importe_linea"]))
)
| "SumPerSKU" >> beam.CombinePerKey(sum)
| "FormatAggregate" >> beam.Map(
lambda kv: {
"dia": kv[0][0],
"sku": kv[0][1],
"ventas_eur": str(kv[1]),
}
)
| "WriteAggregate" >> beam.io.WriteToBigQuery(
table=f"{PROJECT}:{DATASET}.ventas_diarias_sku",
schema="dia:DATE,sku:STRING,ventas_eur:NUMERIC",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE,
)
)
# 6) QUARANTINE: the rejects from both steps, together, to Cloud Storage
(
(parsed[ParseCSVLine.REJECTED_OUTPUT],
validated[ValidateOrder.REJECTED_OUTPUT])
| "MergeRejects" >> beam.Flatten()
| "SerialiseRejects" >> beam.Map(
lambda r: f'{r["reason"]}\t{r["line"]}'
)
| "WriteQuarantine" >> beam.io.WriteToText(
quarantine_path, file_name_suffix=".tsv"
)
)
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--date", required=True, help="Date to process, yyyy-MM-dd")
known, rest = parser.parse_known_args()
options = PipelineOptions(rest)
# Needed so that the workers install the file's dependencies
options.view_as(SetupOptions).save_main_session = True
with beam.Pipeline(options=options) as p:
build_pipeline(p, known.date)
if __name__ == "__main__":
logging.getLogger().setLevel(logging.INFO)
main()Five things worth commenting on:
with beam.Pipeline(...) as p: on leaving thewithblock, Beam callsrun()and waits. Without thewith, you have to callp.run().wait_until_finish()explicitly.WRITE_APPENDversusWRITE_TRUNCATE: the detail is appended (we want accumulated history); the aggregate is rewritten in full (we want the current snapshot). Choosing wrong here silently duplicates data.additional_bq_parameterscreates the table already partitioned and clustered if it did not exist. It is the way not to lose what we learned in 04-01 when the table is created by the pipeline.beam.Flatten()merges severalPCollectionof the same type into one. It is the union of branches in the graph, the equivalent of aUNION ALL.save_main_session=Trueserialises the module's global scope for the workers. Without this, a pipeline that works locally fails on Dataflow withNameErrorabout the constants or the imports. It is beginner mistake number one.
- Running locally with
DirectRunner
DirectRunnerBefore spending a cent, you test locally. The DirectRunner runs the pipeline on your machine, with a subset of the data.
# Local test data
mkdir -p ./pruebas && cat > ./pruebas/pedidos-test.csv <<'EOF'
pedido_id,linea_num,fecha_pedido,sku,cantidad,precio_unitario,descuento_linea,estado
PED-2026-0042,1,2026-03-14,MOCH-40L-AZ,1,89,90,0,confirmado
PED-2026-0042,2,2026-03-14,FRON-300L,2,34.50,5.00,confirmado
PED-2026-0043,1,2026-03-14,CRAM-12P,-1,120.00,0,confirmado
PED-2026-0044,1,fecha-mala,TIEN-2P,1,240.00,0,confirmado
EOF
python pipeline_pedidos.py \
--date 2026-03-14 \
--runner DirectRunnerThat test file is dirty on purpose, and that is how any pipeline's test set should be:
- The first line has nine fields because
89,90slips in an extra comma. The parser will reject it with "expected 8 columns". It is exactly the real failure of a badly configured ERP. - The third has a negative quantity: the validator rejects it.
- The fourth has an invalid date: the validator rejects it.
- Only the second gets through.
The DirectRunner is deliberately strict: it checks the immutability of the elements, it serialises and deserialises between steps, and it shuffles the data out of order on purpose. If your pipeline depends on arrival order or modifies an object in place, the DirectRunner detects it and fails, whereas in the cloud it would produce incorrect results intermittently. Being slow and fussy is the feature, not a defect.
Good practice: as well as testing the whole pipeline, test the transformations with unittest and Beam's utilities:
import unittest
import apache_beam as beam
from apache_beam.testing.test_pipeline import TestPipeline
from apache_beam.testing.util import assert_that, equal_to
class TestValidateOrder(unittest.TestCase):
def test_rejects_negative_quantity(self):
input_rows = [{
"pedido_id": "PED-1", "linea_num": "1", "fecha_pedido": "2026-03-14",
"sku": "CRAM-12P", "cantidad": "-1",
"precio_unitario": "120.00", "descuento_linea": "0", "estado": "confirmado",
}]
with TestPipeline() as p:
outputs = (
p | beam.Create(input_rows)
| beam.ParDo(ValidateOrder()).with_outputs(
ValidateOrder.REJECTED_OUTPUT, main="clean")
)
assert_that(outputs.clean, equal_to([]), label="no valid rows")beam.Create(...) builds a PCollection from a Python list: it is the way to inject test data. assert_that with equal_to compares the contents regardless of order, which is the right thing to do in a distributed system.
- Managed execution with
DataflowRunner
DataflowRunnerWith the tests green, on to the real service.
# A dedicated bucket for Dataflow artefacts (do not mix it with the catalogue)
gcloud storage buckets create gs://alpinashop-dataflow \
--project=alpinashop-datos --location=europe-west1 \
--uniform-bucket-level-access
# Service account dedicated to the pipeline, with least privilege (03-04)
gcloud iam service-accounts create sa-dataflow-pedidos \
--project=alpinashop-datos \
--display-name="Dataflow order pipelines"
SA="[email protected]"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/dataflow.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-dataflow \
--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 launch:
python pipeline_pedidos.py \
--date 2026-03-14 \
--runner DataflowRunner \
--project alpinashop-datos \
--region europe-west1 \
--temp_location gs://alpinashop-dataflow/temp \
--staging_location gs://alpinashop-dataflow/staging \
--service_account_email "$SA" \
--subnetwork regions/europe-west1/subnetworks/sn-datos-euw1 \
--no_use_public_ips \
--machine_type n2-standard-2 \
--max_num_workers 10 \
--job_name alpinashop-pedidos-20260314 \
--labels entorno=produccion,equipo=datos,centro-coste=analiticaOption by option, because each one has consequences:
| Option | What it does and why |
|---|---|
--region europe-west1 |
Where the workers run. It must match the dataset's location and the bucket's, or you will pay cross-region egress and add latency |
--temp_location |
Temporary files and intermediate shuffle. Mandatory |
--staging_location |
Where the pipeline code and its dependencies are uploaded |
--service_account_email |
The workers' identity. Without this the default Compute service account is used, which is usually project Editor: a direct breach of the least privilege of 03-04 |
--subnetwork |
The workers start up inside alpinashop-vpc, in the data subnet. Not on a default network |
--no_use_public_ips |
No public IP: they reach the internet through the Cloud NAT you already configured in 03-01, and reach Google's APIs through Private Google Access |
--machine_type |
The worker's VM type |
--max_num_workers |
Autoscaling ceiling: the safety net against a runaway invoice |
--job_name |
Visible name. Including the date helps locate reprocessing runs |
--labels |
Billing labels, the same scheme as the rest of AlpinaShop |
Checking the status:
gcloud dataflow jobs list --region=europe-west1 \
--format="table(id, name, type, state, createTime)" --limit=5
# Detail and metrics of a specific job
gcloud dataflow jobs describe JOB_ID --region=europe-west1
gcloud dataflow metrics list JOB_ID --region=europe-west1 \
--format="table(name.name, scalar)" --filter="name.name~rows_"That last command returns the rows_ok and rows_rejected counters we instrumented. That is the return on having used Beam metrics instead of print.
- Time in streaming: event versus processing
Here the hard part begins, and also the part that makes learning Beam worthwhile instead of improvising.
In batch, time is simple: you have the file, it has an end, you process and you finish. In streaming there is no end, and a distinction appears that changes everything:
- Event time: when the fact happened in the real world. The customer pressed "Buy" at 10:00:00.
- Processing time: when the data reaches your pipeline. 10:00:03, or 10:07:12 if the phone was in a tunnel, or 11:30 if the app kept the event locally until it got signal back.
With AlpinaShop it is very concrete. A customer on the Barcelona metro browses the catalogue, adds a backpack to the shopping cart at 18:42, loses signal, and the app sends the accumulated events at 18:51.
flowchart LR
subgraph Real["Event time (real world)"]
E1["18:41 ver_producto"]
E2["18:42 anadir_carrito"]
E3["18:44 iniciar_pago"]
end
subgraph Proceso["Processing time (arrival at the pipeline)"]
P1["18:51 all three at once"]
end
E1 --> P1
E2 --> P1
E3 --> P1
Now the business question: how many products were viewed between 18:40 and 18:45? If you group by processing time, the answer is zero, and it is false. If you group by event time, the answer includes this customer, and it is the correct one.
Beam groups by event time by default. That is its most important design decision and the reason its numbers agree with those of the next day's batch report. A system that groups by processing time gives results that depend on the network and are never reproducible: reprocessing the same day gives a different result.
- Watermarks, windows, triggers and late data
If you wait by event time, the inevitable question arises: when do you stop waiting? An 18:42 event could arrive tomorrow. Do you close the window or wait forever?
Beam's answer is four mechanisms that combine.
Windows (Windowing)
They chop the unbounded PCollection into finite chunks over which aggregation is possible.
| Type | Definition | Use at AlpinaShop |
|---|---|---|
| Fixed (tumbling) | Contiguous, non-overlapping intervals | Orders per hour for the campaign dashboard |
| Sliding | Overlapping intervals | Moving average of visits: a 30-minute window every 5 minutes |
| Session | Events separated by less than a given gap are grouped | Real browsing sessions: everything a user does until 30 minutes of inactivity |
| Global | A single infinite window | Only with explicit triggers |
from apache_beam import window
# Fixed 1-hour window: for the order counter
hourly = events | "HourWindow" >> beam.WindowInto(
window.FixedWindows(60 * 60)
)
# Sliding window: 30-minute moving average, updated every 5
moving_average = events | "SlidingWindow" >> beam.WindowInto(
window.SlidingWindows(size=30 * 60, period=5 * 60)
)
# Session window: groups a user's activity
sessions = events | "SessionWindow" >> beam.WindowInto(
window.Sessions(gap_size=30 * 60)
)The session window is particularly elegant and has no simple equivalent in SQL: its duration is not fixed in advance, it is determined by the data itself. It is exactly the definition of "browsing session" that the visitas table from 04-01 needs.
Watermark
The watermark is the system's estimate of "I am no longer expecting events earlier than this instant". Dataflow computes it automatically by observing the timestamps of the data coming in: if events from 18:50 onwards have been arriving from Pub/Sub for ten minutes, the watermark advances past 18:45 and the earlier windows can be closed.
It is not a guarantee, it is a heuristic. Something can always arrive afterwards. That is why the other two mechanisms exist.
Triggers
The trigger decides when to emit a window's result. By default, when the watermark passes: one result per window, when it is considered complete.
But the autumn campaign dashboard cannot wait an hour to see the first number. With a trigger, partial results are emitted:
from apache_beam.transforms.trigger import (
AfterWatermark, AfterProcessingTime, AccumulationMode
)
orders_per_hour = (
events
| "Window" >> beam.WindowInto(
window.FixedWindows(60 * 60),
trigger=AfterWatermark(
early=AfterProcessingTime(60), # partial update every minute
late=AfterProcessingTime(10 * 60), # corrections every 10 min if data is late
),
allowed_lateness=2 * 60 * 60, # we accept up to 2 h of delay
accumulation_mode=AccumulationMode.ACCUMULATING,
)
| "Count" >> beam.CombinePerKey(sum)
)A practical reading of this configuration for AlpinaShop:
- Every minute a provisional result is emitted: the dashboard moves and people see that the system is alive.
- When the watermark passes, the result considered final is emitted.
- For two more hours (
allowed_lateness), if events from that hour arrive — the customer on the metro — a correction is emitted every ten minutes. - After those two hours, whatever arrives is discarded. That discard is an explicit business decision, not an accident, and it has to be instrumented with a counter so you know how much is being lost.
AccumulationMode.ACCUMULATING means each emission contains the accumulated total for the window, so the destination must overwrite. The alternative, DISCARDING, emits only what is new since the last emission, and the destination must add it up. Confusing them produces doubled or halved figures, and it is a very hard mistake to spot on a dashboard.
Late data
Everything that arrives after the watermark is late data. With allowed_lateness you decide how long you accept it. The rule is simple: the longer you wait, the more correct the numbers are and the more resources (in-memory state) the pipeline consumes. Two hours for a European e-commerce business is a reasonable value; two days would be extremely expensive and would not change any decision.
- The
pedidos-nuevos streaming pipeline
pedidos-nuevos streaming pipelineIn the next lesson, 04-04, we will create the Pub/Sub topic pedidos-nuevos, where the Flask application will publish a JSON message every time an order is confirmed. Let us anticipate the consumer, because it is Dataflow's canonical use case.
"""
Streaming pipeline: pedidos-nuevos (Pub/Sub) -> BigQuery.
It runs continuously, with no end.
"""
import json
import logging
from datetime import datetime
import apache_beam as beam
from apache_beam import window
from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions
from apache_beam.transforms.trigger import AfterWatermark, AfterProcessingTime, AccumulationMode
PROJECT = "alpinashop-datos"
SUBSCRIPTION = f"projects/{PROJECT}/subscriptions/sub-analitica"
class DecodeOrder(beam.DoFn):
REJECTED_OUTPUT = "rejected"
def process(self, message, moment=beam.DoFn.TimestampParam):
try:
data = json.loads(message.data.decode("utf-8"))
except Exception as exc:
yield beam.pvalue.TaggedOutput(
self.REJECTED_OUTPUT,
{"payload": str(message.data[:500]), "reason": f"invalid json: {exc}"},
)
return
# The Pub/Sub message attributes travel separately from the body
source = message.attributes.get("origen", "desconocido")
yield {
"pedido_id": data["pedido_id"],
"fecha_pedido": data["fecha"][:10],
"momento_pedido": data["fecha"],
"cliente_id": data.get("cliente_id"),
"canal": source,
"estado": "confirmado",
"total_pedido": str(data["total"]),
"momento_ingesta": datetime.utcnow().isoformat(),
}
def main():
options = PipelineOptions(streaming=True, save_main_session=True)
options.view_as(StandardOptions).streaming = True
with beam.Pipeline(options=options) as p:
messages = (
p
| "ReadPubSub" >> beam.io.ReadFromPubSub(
subscription=SUBSCRIPTION,
with_attributes=True,
timestamp_attribute="momento_evento", # <-- the key line
)
)
decoded = (
messages
| "Decode" >> beam.ParDo(DecodeOrder()).with_outputs(
DecodeOrder.REJECTED_OUTPUT, main="valid")
)
# A) Real-time detail, row by row
(
decoded.valid
| "WriteOrders" >> beam.io.WriteToBigQuery(
table=f"{PROJECT}:alpinashop_analitica.pedidos_streaming",
schema="pedido_id:STRING,fecha_pedido:DATE,momento_pedido:TIMESTAMP,"
"cliente_id:STRING,canal:STRING,estado:STRING,"
"total_pedido:NUMERIC,momento_ingesta:TIMESTAMP",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
method="STORAGE_WRITE_API",
)
)
# B) Aggregate per hour and channel, for the campaign dashboard
(
decoded.valid
| "HourWindow" >> beam.WindowInto(
window.FixedWindows(3600),
trigger=AfterWatermark(early=AfterProcessingTime(60)),
allowed_lateness=7200,
accumulation_mode=AccumulationMode.ACCUMULATING,
)
| "KeyByChannel" >> beam.Map(lambda d: (d["canal"], float(d["total_pedido"])))
| "SumPerChannel" >> beam.CombinePerKey(sum)
| "Format" >> beam.Map(
lambda kv, w=beam.DoFn.WindowParam: {
"hora_inicio": w.start.to_utc_datetime().isoformat(),
"canal": kv[0],
"ventas_eur": round(kv[1], 2),
}
)
| "WriteAggregate" >> beam.io.WriteToBigQuery(
table=f"{PROJECT}:alpinashop_analitica.ventas_por_hora",
schema="hora_inicio:TIMESTAMP,canal:STRING,ventas_eur:FLOAT",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
method="STORAGE_WRITE_API",
)
)
if __name__ == "__main__":
logging.getLogger().setLevel(logging.INFO)
main()Three decisions that have to be understood:
timestamp_attribute="momento_evento": it tells Beam to use the message attribute as the event time, instead of the moment of publication in Pub/Sub. Without this, a message held for nine minutes on the phone would be counted in the wrong hour. It is the line that makes the whole of section 9 actually work, and it is constantly forgotten.method="STORAGE_WRITE_API": it uses the modern write API, cheaper and with exactly-once semantics, instead of the classic streaming inserts.- It reads from a subscription, not from the topic. With a subscription, if the pipeline stops, the messages pile up and are recovered on start-up. Reading from the topic directly makes Dataflow create an ephemeral subscription and everything published while the pipeline is stopped is lost.
- Dataflow templates: the practical route
Everything above is powerful and it is also work. For common tasks there is something much simpler: templates.
A template is an already compiled and parameterised pipeline that is launched without writing or compiling any code. Google publishes dozens of them.
| Template | What it does |
|---|---|
| Pub/Sub Subscription to BigQuery | Reads JSON from a subscription and inserts it into a table |
| Cloud Storage Text to BigQuery | Loads files with a JavaScript transformation function |
| JDBC to BigQuery | Dumps a relational database |
| BigQuery to Cloud Storage (Parquet) | Exports |
| Datastream to BigQuery | Applies CDC (we will see it in 04-05) |
| Bulk Compress/Decompress | Utilities over the bucket |
Launching the Pub/Sub to BigQuery one for AlpinaShop is a single line:
gcloud dataflow jobs run alpinashop-pedidos-a-bq \
--gcs-location gs://dataflow-templates-europe-west1/latest/PubSub_Subscription_to_BigQuery \
--region europe-west1 \
--service-account-email "[email protected]" \
--subnetwork regions/europe-west1/subnetworks/sn-datos-euw1 \
--disable-public-ips \
--max-workers 5 \
--parameters \
inputSubscription=projects/alpinashop-datos/subscriptions/sub-analitica,\
outputTableSpec=alpinashop-datos:alpinashop_analitica.pedidos_streaming,\
outputDeadletterTable=alpinashop-datos:alpinashop_analitica.pedidos_streaming_erroresNotice outputDeadletterTable: messages that do not fit the schema go to an error table instead of blocking the pipeline. It is the same quarantine pattern we programmed by hand, already solved.
There are two flavours:
- Classic templates: the graph is compiled when the template is created; parameters can only be runtime values (
ValueProvider). - Flex Templates: the pipeline is packaged as a container image in Artifact Registry and the graph is built at launch time. They are more flexible and are the recommended option today for your own templates.
Creating a flex template of AlpinaShop's batch pipeline, so that Composer or Cloud Scheduler can invoke it in 04-06:
# 1) Build the image and publish the template
gcloud dataflow flex-template build \
gs://alpinashop-dataflow/plantillas/pedidos-lote.json \
--image-gcr-path europe-west1-docker.pkg.dev/alpinashop-prod/alpinashop/dataflow-pedidos:1.0.0 \
--sdk-language PYTHON \
--flex-template-base-image PYTHON3 \
--py-path . \
--env FLEX_TEMPLATE_PYTHON_PY_FILE=pipeline_pedidos.py \
--env FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE=requirements.txt
# 2) Run it with parameters
gcloud dataflow flex-template run "pedidos-$(date +%Y%m%d-%H%M%S)" \
--template-file-gcs-location gs://alpinashop-dataflow/plantillas/pedidos-lote.json \
--region europe-west1 \
--service-account-email "[email protected]" \
--parameters date=2026-03-14The image goes to the same Artifact Registry the catalogue already uses (europe-west1-docker.pkg.dev/alpinashop-prod/alpinashop/), which keeps a single inventory of artefacts, consistent with what we will see in module 6.
AlpinaShop's criterion: for "Pub/Sub to BigQuery" with no transformation, Google's template, without writing a line. For the batch pipeline with its own business validation, Beam code packaged as a flex template. Write code only when it adds logic the template does not have.
- Autoscaling and Dataflow Prime
Dataflow adjusts the number of workers during execution. In batch, it looks at the pending work; in streaming, it looks at the queue backlog and CPU usage.
--num_workers 2 # how many it starts with
--max_num_workers 20 # hard ceiling: the cost control
--autoscaling_algorithm THROUGHPUT_BASED # the default in streamingSetting a low max_num_workers does not always save money: a pipeline that takes ten hours with two workers can cost the same as one that takes an hour with twenty, because you are billed per worker per hour. What the ceiling does prevent is the surprise of a job with a pathological loop consuming a hundred machines all night.
Dataflow Prime is the evolution of the service, with three practical differences:
| Aspect | Classic Dataflow | Dataflow Prime |
|---|---|---|
| Resources | You choose the machine type for the whole pipeline | Automatic vertical memory adjustment per step |
| Scaling | Horizontal | Horizontal + vertical |
| Billing | Per vCPU, memory and disk per hour | Per Data Compute Units (DCU) |
| Diagnostics | Metrics | Automatic bottleneck recommendations |
| When to use it | Stable, well-sized pipelines | Pipelines with steps of very uneven consumption |
Prime's real advantage shows up when one step of the pipeline needs a lot of memory and the others do not: in the classic model you size every machine for the most demanding step and waste resources the rest of the time. It is enabled with --dataflow_service_options=enable_prime.
For AlpinaShop, with modest and predictable pipelines, the classic model with n2-standard-2 is enough and easier to reason about on the invoice. Prime is noted down for when the volume justifies it.
- Monitoring, parallelism and data skew
The Dataflow console shows the execution graph with each step and, for each one, elements processed, CPU time and status. It is the best data-debugging interface on the platform, and that is why we insist on naming the steps.
Metrics you have to watch:
| Metric | What it indicates | Alarm threshold |
|---|---|---|
| System lag | Seconds the oldest unprocessed element has been waiting | In streaming, if it grows steadily, the pipeline cannot keep up |
| Data freshness | Age of the most recent data already emitted | It must stay stable |
| Elements per second per step | Where the bottleneck is | The slowest step sets the pace |
| vCPU usage | Whether scaling is helping | High usage with growing lag: there is a structural limit |
| Current workers | Autoscaling behaviour | Pinned to the maximum: raise the ceiling or fix the pipeline |
The most frequent and hardest problem to diagnose is data skew: one key concentrates an enormous proportion of the elements. At AlpinaShop it is very easy for it to happen: if you group visits by sku and the flagship backpack accumulates 40 % of the traffic, a single worker will process 40 % of the work while the other nineteen wait. The symptom is unmistakable: scaling does not improve anything and there is one step with one worker at 100 % and the rest idle.
Three remedies, in order of preference:
# 1) THE BEST: use CombinePerKey instead of GroupByKey.
# Partial aggregation happens on each worker before the shuffle.
sales = lines | "Sum" >> beam.CombinePerKey(sum)
# 2) If you genuinely need GroupByKey: add salt to the key
import random
def salt(element, n=20):
key, value = element
return (f"{key}#{random.randint(0, n - 1)}", value)
result = (
lines
| "Salt" >> beam.Map(salt)
| "PartialGroup" >> beam.CombinePerKey(sum)
| "RemoveSalt" >> beam.Map(lambda kv: (kv[0].split("#")[0], kv[1]))
| "FinalGroup" >> beam.CombinePerKey(sum)
)
# 3) If one side of the JOIN is small (the product catalogue):
# use a side input instead of CoGroupByKey, and avoid the shuffle
products = p | "ReadProducts" >> beam.io.ReadFromBigQuery(query=SQL_PRODUCTS)
catalogue = beam.pvalue.AsDict(products | beam.Map(lambda r: (r["sku"], r)))
enriched = lines | "Enrich" >> beam.Map(
lambda line, cat: {**line, "categoria": cat.get(line["sku"], {}).get("categoria")},
cat=catalogue,
)Technique 2, the salt, deserves an explanation: by adding a random suffix to the key, the flagship backpack becomes twenty different keys that are shared out across twenty workers; then the salt is removed and a second aggregation is done over twenty values, which is trivial. It only works with associative operations, but it covers almost every aggregation case.
Technique 3, the side input, is the one you will use most at AlpinaShop: the product catalogue is a few thousand rows and fits in each worker's memory. Broadcasting it avoids the JOIN shuffle entirely. Watch out for the limit: if the side input does not fit in memory, the pipeline degrades enormously.
- Cost: what exactly you pay for
Dataflow does not bill per pipeline or per unit of data processed: it bills the resources consumed by the workers.
| Item | How it is measured | Order of magnitude (verify in the official documentation) |
|---|---|---|
| vCPU | Per vCPU per hour | ~$0.05-0.07 (batch); somewhat more in streaming |
| Memory | Per GB per hour | ~$0.003-0.004 |
| Persistent disk | Per GB per hour | ~$0.00005 standard |
| Shuffle (batch) | Per GB processed by the shuffle service | ~$0.011/GB |
| Streaming Engine | Per GB of streaming data processed | ~$0.018/GB |
| Dataflow Prime | Per DCU | Unified model |
A realistic example for AlpinaShop: the nightly batch pipeline, with 4 n2-standard-2 workers (2 vCPU, 8 GB) for 20 minutes, comes out at a few cents. The streaming pipeline, by contrast, runs 24×7: 2 permanent workers are around 1,440 vCPU-hours a month, on the order of €80-100 monthly. That is the number that surprises everybody, and it is the reason for the warning at the start of the lesson.
Concrete tips for not overpaying:
- Enable Streaming Engine and Shuffle Service (
--enable_streaming_engine,--experiments=shuffle_mode=service). They move the shuffle and the state outside the workers, allowing smaller machines and much smaller disks. It almost always pays off. - Reduce the disk. With Streaming Engine,
--disk_size_gb=30is enough; the default value is much larger and it is billed per hour. - Always set
--max_num_workers. It is the handbrake. - Use Spot VMs in batch that tolerates interruptions:
--flexrs_goal=COST_OPTIMIZEDdelays the start by up to 6 hours in exchange for a notable discount. Perfect for the nightly dump; unacceptable for streaming. - Do you really need streaming? A micro-batch every 15 minutes with the Cloud Storage to BigQuery template costs a fraction of a permanent pipeline. If the business tolerates 15 minutes of delay, the answer is no.
- Watch out for forgotten streaming pipelines. A test job nobody stopped is the most common phantom line item on a data invoice. You will spot it with
gcloud dataflow jobs list --status=active.
- When Dataflow is not the answer
Being honest about the limits is part of knowing how to use a tool.
| Situation | Better option | Why |
|---|---|---|
| A transformation expressible in SQL over data already in BigQuery | BigQuery (04-01) | Do not move data in order to transform it; use INSERT ... SELECT or a materialized view |
| You already have Spark code, or the team knows Spark rather than Beam | Dataproc (04-03) | Rewriting to Beam is a cost with no clear return |
| A simple file load into BigQuery, with no logic | bq load |
It is free; Dataflow would cost money to do the same thing |
| Integrating external sources with a team that does not program | Data Fusion (04-05) | Visual interface, ready-made connectors |
| Deciding the order in which several processes run | Composer or Workflows (04-06) | Dataflow runs a pipeline; it does not orchestrate the others |
| Reacting to a one-off event with little logic | Cloud Functions (06-03) | A whole pipeline to process one file is disproportionate |
| Copying a database with change capture | Datastream (04-05) | Managed CDC, no code |
The most common confusion is the one in the penultimate row. Dataflow is not an orchestrator. It can run an extremely complex pipeline, but it does not know how to do "first export Cloud SQL, then load into BigQuery, then launch this pipeline, and if anything fails alert Marta". That is exactly what 04-06 solves.
Common Mistakes and Tips
Forgetting save_main_session=True. The pipeline works locally and fails on Dataflow with NameError: name 'PROJECT' is not defined. The workers do not receive the module's global scope unless you tell them to.
Not pinning the SDK version. The workers use the version you launched the pipeline with. A pip install apache-beam[gcp] with no version means the same code works today and fails tomorrow. Pin the version in requirements.txt.
Not naming the steps. Without "Name" >>, the console graph is illegible and the metrics say nothing. On top of that, changing a step's name prevents you from updating a running streaming pipeline (--update), because Beam cannot map the old step's state onto the new one.
Using GroupByKey where CombinePerKey would do. It is the difference between a pipeline that scales and one that falls over on a hot key.
Reading from a topic instead of a subscription. With a topic, Dataflow creates a temporary subscription and everything published while the pipeline is stopped is lost. With a subscription, the messages pile up and are recovered.
Forgetting timestamp_attribute. The windows are computed over the publication moment instead of the event moment, and the numbers do not agree with the batch process. It is subtle and serious.
Leaving a test streaming pipeline switched on. It bills 24×7. Put labels on them and review gcloud dataflow jobs list --status=active regularly.
Not using your own service account. Without --service_account_email, the workers use the default Compute Engine account, which in many projects is Editor. A pipeline with Editor permissions over production is an unnecessary risk.
Tip: always write the error branch. A pipeline with no reject output is not a production pipeline. If you cannot see what was discarded and why, you cannot trust the numbers.
Tip: --update to modify a streaming pipeline. It lets you replace the code while preserving the in-flight state, instead of draining it and starting from scratch. It requires graph compatibility, which is another reason not to rename steps casually.
Tip: drain, do not cancel. gcloud dataflow jobs drain stops reading new input and finishes what it has in flight. cancel kills the job and can lose data in process.
Exercises
Exercise 1: a customer review pipeline with quarantine
Write a batch pipeline in Beam that reads gs://alpinashop-catalogo/exportaciones/2026/03/opiniones-*.csv with the columns opinion_id,sku,fecha,puntuacion,texto,pais, and that:
- rejects rows whose
puntuacionis not an integer between 1 and 5, or whosepaisis not a two-letter code; - normalises the
skuto uppercase and trimstextoto 500 characters; - writes the valid ones into
alpinashop-datos:alpinashop_analitica.opiniones; - writes the rejected ones, with their reason, into
gs://alpinashop-catalogo/cuarentena/opiniones/; - keeps a counter of valid and rejected rows visible in Dataflow.
Test it with DirectRunner and a local file with at least two faulty rows.
Exercise 2: windows and triggers for the campaign dashboard
Management wants an autumn campaign dashboard with sales per hour and per country. Requirements: group by event time, show a partial update every 30 seconds so the dashboard looks alive, accept events with up to 90 minutes of delay, emitting corrections, and have each emission contain the accumulated total for the hour. Write only the WindowInto fragment and the aggregation, and explain what exactly the destination writes and why the chosen accumulation mode forces a specific write_disposition.
Exercise 3: diagnosing a pipeline that does not scale
The visits streaming pipeline has been running for three days. Since yesterday, the system lag has gone from 4 seconds to 22 minutes and it keeps growing. Autoscaling has reached 20 workers (its maximum), but the average vCPU usage across them is 18 %. In the graph, the GroupBySKU step shows one worker with 6 hours of CPU time and the rest with less than 10 minutes. Yesterday, marketing launched a campaign for the MOCH-40L-AZ backpack that has multiplied its visits twelvefold.
Diagnose the cause, explain why raising --max_num_workers to 50 would not fix anything, and propose two code solutions with their practical difference.
Solutions
Solution 1
import csv, io, logging, re
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions, SetupOptions
PROJECT, DATASET = "alpinashop-datos", "alpinashop_analitica"
BUCKET = "alpinashop-catalogo"
RE_COUNTRY = re.compile(r"^[A-Z]{2}$")
class ProcessReview(beam.DoFn):
REJECTS = "rejects"
def __init__(self):
self.ok = beam.metrics.Metrics.counter("reviews", "valid")
self.ko = beam.metrics.Metrics.counter("reviews", "rejected")
def process(self, line):
try:
c = next(csv.reader(io.StringIO(line)))
except Exception as exc:
self.ko.inc()
yield beam.pvalue.TaggedOutput(self.REJECTS,
{"line": line, "reason": f"unreadable csv: {exc}"})
return
if len(c) != 6:
self.ko.inc()
yield beam.pvalue.TaggedOutput(self.REJECTS,
{"line": line, "reason": f"expected 6 columns, found {len(c)}"})
return
opinion_id, sku, fecha, score, texto, pais = [x.strip() for x in c]
reasons = []
try:
puntuacion = int(score)
if not 1 <= puntuacion <= 5:
reasons.append("score outside the 1-5 range")
except ValueError:
reasons.append("non-numeric score")
puntuacion = None
pais = pais.upper()
if not RE_COUNTRY.match(pais):
reasons.append(f"invalid country: {pais}")
if not sku:
reasons.append("empty sku")
if reasons:
self.ko.inc()
yield beam.pvalue.TaggedOutput(self.REJECTS,
{"line": line, "reason": "; ".join(reasons)})
return
self.ok.inc()
yield {
"opinion_id": opinion_id,
"sku": sku.upper(),
"fecha": fecha,
"puntuacion": puntuacion,
"texto": texto[:500],
"pais": pais,
}
def main():
options = PipelineOptions()
options.view_as(SetupOptions).save_main_session = True
with beam.Pipeline(options=options) as p:
outputs = (
p
| "Read" >> beam.io.ReadFromText(
f"gs://{BUCKET}/exportaciones/2026/03/opiniones-*.csv",
skip_header_lines=1)
| "Process" >> beam.ParDo(ProcessReview()).with_outputs(
ProcessReview.REJECTS, main="valid")
)
(outputs.valid
| "WriteBQ" >> beam.io.WriteToBigQuery(
table=f"{PROJECT}:{DATASET}.opiniones",
schema="opinion_id:STRING,sku:STRING,fecha:DATE,"
"puntuacion:INTEGER,texto:STRING,pais:STRING",
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
additional_bq_parameters={
"timePartitioning": {"type": "DAY", "field": "fecha"},
"clustering": {"fields": ["sku"]},
}))
(outputs[ProcessReview.REJECTS]
| "Serialise" >> beam.Map(lambda r: f'{r["reason"]}\t{r["line"]}')
| "WriteQuarantine" >> beam.io.WriteToText(
f"gs://{BUCKET}/cuarentena/opiniones/rechazos",
file_name_suffix=".tsv"))
if __name__ == "__main__":
logging.getLogger().setLevel(logging.INFO)
main()Test file with deliberate defects:
opinion_id,sku,fecha,puntuacion,texto,pais OPI-0001,moch-40l-az,2026-03-10,5,Very comfortable for long treks,es OPI-0002,FRON-300L,2026-03-10,9,Impossible score,FR OPI-0003,CRAM-12P,2026-03-11,four,Non-numeric score,ES OPI-0004,TIEN-2P,2026-03-11,4,Malformed country,SPAIN
The first one gets through (the sku is normalised to uppercase and es to ES); the next three go to quarantine with different reasons. Expected metrics: valid=1, rejected=3.
Solution 2
from apache_beam import window
from apache_beam.transforms.trigger import (
AfterWatermark, AfterProcessingTime, AccumulationMode
)
sales_hour_country = (
events
| "CampaignWindow" >> beam.WindowInto(
window.FixedWindows(3600), # 1 hour, EVENT time
trigger=AfterWatermark(
early=AfterProcessingTime(30), # update every 30 s
late=AfterProcessingTime(300), # corrections every 5 min
),
allowed_lateness=90 * 60, # 90 minutes of delay
accumulation_mode=AccumulationMode.ACCUMULATING,
)
| "KeyByCountry" >> beam.Map(lambda d: (d["pais"], float(d["total_pedido"])))
| "SumPerCountry" >> beam.CombinePerKey(sum)
| "Format" >> beam.Map(
lambda kv, w=beam.DoFn.WindowParam: {
"hora_inicio": w.start.to_utc_datetime().isoformat(),
"pais": kv[0],
"ventas_eur": round(kv[1], 2),
})
)What exactly the destination writes. With ACCUMULATING, each emission of a window contains the accumulated total since the start of the window, not the increment. For the 10:00-11:00 window and country ES, the destination will receive a sequence like: €120 (at 10:00:30), €345 (10:01:00), … €4,210 (when the watermark passes), €4,235 (at 11:20, a correction for a late event).
Why it forces a specific write_disposition. If WRITE_APPEND were used on the final table, those six emissions would add up and the dashboard would show more than €9,000 where there are €4,235: the figure would be multiplied. The correct options are:
- Write into an intermediate state table with
WRITE_APPENDand have the dashboard query only the last emission per window and key, for example withQUALIFY ROW_NUMBER() OVER (PARTITION BY hora_inicio, pais ORDER BY momento_emision DESC) = 1. This is the recommended option in streaming, because it preserves the history of corrections and allows auditing. - Or use a destination that supports overwriting by key (a later
MERGEonhora_inicio+pais).
If AccumulationMode.DISCARDING were chosen instead, each emission would carry only the increment and then WRITE_APPEND with a later sum would be the right thing. The combination of accumulation mode plus write disposition has to be decided jointly; getting it wrong produces dashboards with inflated figures that nobody detects until somebody cross-checks the number with accounting.
Solution 3
Diagnosis: data skew caused by a hot key.
The three pieces of evidence all point the same way and reinforce each other:
- 20 workers at 18 % average CPU. If the pipeline were genuinely saturated, the CPU would be high. Low usage with growing lag means the workers are waiting, not working.
- One worker with 6 hours of CPU and the rest with 10 minutes. That is the exact signature of skew: the work is not being shared out.
- The
MOCH-40L-AZcampaign. In aGroupByKeybysku, every event for the flagship backpack is sent through the shuffle to a single worker, because the key determines the destination. The others cannot help it.
Why going up to 50 workers fixes nothing. The unit of parallelism in a grouping by key is the key, not the element. The MOCH-40L-AZ events will all keep going to the same place, with 20 workers or with 500. The 30 new ones would sit idle and bill: the lag would keep growing and the cost would be multiplied by 2.5. Scaling horizontally does not solve a distribution problem.
Solution A — replace GroupByKey with CombinePerKey (the good one, if the operation allows it):
# BEFORE: every event for the hot key travels to one worker
visits_per_sku = (
events
| "KeyBySKU" >> beam.Map(lambda e: (e["sku"], 1))
| "Group" >> beam.GroupByKey()
| "Count" >> beam.Map(lambda kv: (kv[0], len(list(kv[1]))))
)
# AFTER: each worker sums its own share before the shuffle
visits_per_sku = (
events
| "KeyBySKU" >> beam.Map(lambda e: (e["sku"], 1))
| "Count" >> beam.CombinePerKey(sum)
)What travels over the network is partial sums, one per worker per key, instead of millions of individual elements. With 20 workers, the MOCH-40L-AZ worker receives 20 numbers instead of twelve million events. It is one line of code and it usually solves the problem completely.
Solution B — add salt to the key (when you genuinely need GroupByKey, for example to keep the elements):
import random
N_SALT = 50
visits_per_sku = (
events
| "SaltKey" >> beam.Map(
lambda e: (f'{e["sku"]}#{random.randint(0, N_SALT - 1)}', 1))
| "SaltedPartial" >> beam.CombinePerKey(sum)
| "RemoveSalt" >> beam.Map(lambda kv: (kv[0].split("#")[0], kv[1]))
| "FinalTotal" >> beam.CombinePerKey(sum)
)Practical difference between A and B. A is preferable whenever it is possible: it is simpler, cheaper and has no parameters to tune. B adds an extra shuffle and a parameter (N_SALT) that has to be sized — too low does not spread the load, too high creates overhead — but it is the only route when the operation is not associative (for example, if you need the complete list of a session's events to reconstruct a journey) or when the skew persists even with partial aggregation.
An immediate complementary measure: before touching any code, lower --max_num_workers to 5 to stop paying for 15 idle machines while the deployment is prepared, and use --update to replace the pipeline while preserving the in-flight state, without losing the data in flight.
Conclusion
AlpinaShop now has pipes. In this lesson you have seen why a script on a VM is a solution that works until it stops working, and exactly what a managed service adds: automatic parallelism, retries, scaling during execution, reliable write semantics and observability as standard.
You have learned the Apache Beam model — Pipeline, immutable PCollection, PTransform, runner — and its central idea: a batch is a bounded flow and a stream is an unbounded flow, so the same logic serves both. You know the transformations that cover almost every case and, above all, you know why CombinePerKey is almost always better than GroupByKey, and when a ParDo with a DoFn class beats a Map.
You have written AlpinaShop's batch pipeline line by line: it reads the CSV exports from alpinashop-catalogo, parses, validates with real business rules — amounts with commas, negative quantities, broken dates — writes the clean detail into lineas_pedido, computes the sales aggregate per SKU with partial aggregation, and sends everything faulty to a quarantine in the bucket with its reason, without bringing the process down. You have tested it locally with the DirectRunner, which is slow and strict on purpose, using a deliberately dirty file, and then launched it on Dataflow with the sa-dataflow-pedidos service account, inside sn-datos-euw1, with no public IP, with a worker ceiling and billing labels.
You have entered the territory of time, which is what really distinguishes serious data processing: event time versus processing time, with the customer on the metro as the example of why grouping by the latter produces false numbers; fixed, sliding and session windows; watermarks as the estimate of "I am no longer waiting"; triggers to emit partial results without giving up later correction; and allowed_lateness as an explicit business decision about how long you wait for the stragglers. And you have written the streaming pipeline that will consume pedidos-nuevos as soon as it exists, with timestamp_attribute — the line that makes everything above true — and reading from a subscription, not a topic.
You have seen that for the common cases there is a shortcut: Google's templates solve "Pub/Sub to BigQuery" without writing code, and flex templates package your own pipeline as an image in Artifact Registry so that another system can invoke it. You know autoscaling and its limits, Dataflow Prime and when it is worth it, and you can read the graph, the system lag and the data freshness to diagnose problems. You can recognise data skew — the scaling that improves nothing — and correct it with partial aggregation, salted keys or side inputs. And you know the invoice: per vCPU, memory, disk and shuffle, with the big warning that a streaming pipeline is always switched on and costs tens of euros a month even when nothing happens.
There is something still outstanding that you have noticed all through the lesson. The batch pipeline works because the data is already in the bucket, and the streaming one works because somebody publishes to pedidos-nuevos. But that topic does not exist yet, and we have not talked about what happens to the rest of the house: when an order comes in, the warehouse has to prepare it, billing has to issue the invoice, the customer has to receive their confirmation email and analytics has to find out. Today the Flask application would have to call all four, one after the other, and hang if the fourth does not respond.
Before solving that, however, there is a piece of the data ecosystem worth knowing, because many people arrive at Google Cloud with it already in place. In 04-03, Cloud Dataproc, we will look at managed Spark and Hadoop: what they are, why they still matter after twenty years, and what the pattern is that turns an expensive, permanent cluster into an ephemeral one that lives for four minutes, does its job over the data in the bucket and self-destructs. We will create alpinashop-spark, run a PySpark job that works out which products are bought together — the basis of the future recommender in module 5 — and compare honestly when Dataproc is the right choice, when Dataflow is, and when neither of the two is needed.
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
