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

  1. What a data pipeline is and why a script is not enough
  2. Apache Beam: the unified model
  3. The four concepts: Pipeline, PCollection, PTransform, runner
  4. The transformations you will use 90 % of the time
  5. The first batch pipeline, explained line by line
  6. Running locally with DirectRunner
  7. Managed execution with DataflowRunner
  8. Time in streaming: event versus processing
  9. Watermarks, windows, triggers and late data
  10. The pedidos-nuevos streaming pipeline
  11. Dataflow templates: the practical route
  12. Autoscaling and Dataflow Prime
  13. Monitoring, parallelism and data skew
  14. Cost: what exactly you pay for
  15. When Dataflow is not the answer

  1. 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.

  1. 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.

  1. The four concepts: Pipeline, PCollection, PTransform, runner

Pipeline 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:

result = input_collection | "Descriptive name of the step" >> beam.Map(function)

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.

  1. 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.

  1. 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:

  • yield instead of return. A DoFn is a generator: it can emit zero, one or many elements per input. A return with a value does not work the way you expect.
  • TaggedOutput marks 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 exports 89,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:

  1. with beam.Pipeline(...) as p: on leaving the with block, Beam calls run() and waits. Without the with, you have to call p.run().wait_until_finish() explicitly.
  2. WRITE_APPEND versus WRITE_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.
  3. additional_bq_parameters creates 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.
  4. beam.Flatten() merges several PCollection of the same type into one. It is the union of branches in the graph, the equivalent of a UNION ALL.
  5. save_main_session=True serialises the module's global scope for the workers. Without this, a pipeline that works locally fails on Dataflow with NameError about the constants or the imports. It is beginner mistake number one.

  1. Running locally with DirectRunner

Before 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 DirectRunner

That 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,90 slips 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.

  1. Managed execution with DataflowRunner

With 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=analitica

Option 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.

  1. 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.

  1. 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.

  1. The pedidos-nuevos streaming pipeline

In 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.

  1. 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_errores

Notice 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-14

The 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.

  1. 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 streaming

Setting 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.

  1. 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.

  1. 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=30 is 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_OPTIMIZED delays 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.

  1. 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 puntuacion is not an integer between 1 and 5, or whose pais is not a two-letter code;
  • normalises the sku to uppercase and trims texto to 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_APPEND and have the dashboard query only the last emission per window and key, for example with QUALIFY 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 MERGE on hora_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:

  1. 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.
  2. 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.
  3. The MOCH-40L-AZ campaign. In a GroupByKey by sku, 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

Module 2: Core GCP Services

Module 3: Networking and Security

Module 4: Data and Analytics

Module 5: Machine Learning and AI

Module 6: DevOps and Monitoring

Module 7: Advanced GCP Topics

Module 8: Final Project

© Copyright 2026. All rights reserved