GCP Data Pipelines
Moving and transforming data between systems is a different job than storing it. databases.md covers where data lives and bigquery.md already covers the ingestion paths that land data directly into BigQuery (Storage Write API, batch load jobs). This guide covers the layer that usually sits in front of those paths, or moves data somewhere that isn't BigQuery at all: Dataflow (Apache Beam, unified batch + streaming transforms), Dataproc (managed Hadoop/Spark), Cloud Composer and Workflows (orchestration), and Data Fusion (visual/no-code ETL).
Pipeline Service Map
| Use case | AWS | GCP |
|---|---|---|
| Unified batch + streaming transform (Apache Beam) | Kinesis Data Analytics / Glue Streaming | Dataflow |
| Managed Hadoop / Spark | EMR | Dataproc |
| Complex DAG orchestration across many systems | MWAA (managed Airflow) | Cloud Composer |
| Lightweight serverless step orchestration | Step Functions | Workflows |
| Cron trigger layer | EventBridge Scheduler | Cloud Scheduler |
| No-code / visual ETL builder | Glue Studio | Data Fusion |
| Land data straight into a warehouse | Kinesis Firehose | Streaming inserts (see bigquery.md) |
Dataflow — Apache Beam, One Model for Batch and Streaming
Apache Beam is a programming model, not a service: you write a pipeline as a graph of PTransforms operating on PCollections (an unbounded or bounded set of elements), and a runner executes that graph. Dataflow is Google's fully managed, serverless runner for Beam — the same Beam SDK also runs on Spark or Flink runners elsewhere, but Dataflow is the GCP-native, no-clusters-to-manage option.
The defining feature of the model: the same transform code runs unchanged against a bounded batch source or an unbounded streaming source. Only the source, the windowing, and the sink change — the business logic in between doesn't know or care which mode it's running in.
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
with beam.Pipeline(options=PipelineOptions()) as p:
(p
| "Read" >> beam.io.ReadFromText("gs://my-bucket/orders/*.csv")
| "Parse" >> beam.Map(parse_csv_row)
| "SumByUser" >> beam.CombinePerKey(sum)
| "Write" >> beam.io.WriteToBigQuery(
"project:dataset.user_totals",
write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE))
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions
options = PipelineOptions(streaming=True)
with beam.Pipeline(options=options) as p:
(p
| "Read" >> beam.io.ReadFromPubSub(topic="projects/my-project/topics/orders")
| "Parse" >> beam.Map(parse_json_message)
| "Window" >> beam.WindowInto(beam.window.FixedWindows(60))
| "SumByUser" >> beam.CombinePerKey(sum)
| "Write" >> beam.io.WriteToBigQuery(
"project:dataset.user_totals_live",
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND))
Parse and SumByUser are identical in both pipelines. The only differences are the source (ReadFromText vs ReadFromPubSub), the addition of a Window stage (a batch job has a natural end, so it doesn't need one), and the write disposition (TRUNCATE for a full batch recompute, APPEND for a continuous stream of results). That's the entire pitch of the unified model: build and test your transform logic once, in batch, against a static sample file — then point the same code at a live topic for production.
Autoscaling Workers
Dataflow adjusts worker count automatically based on backlog (streaming) or estimated remaining work (batch), scaling from a couple of workers up to whatever --max-workers allows. A burst of Pub/Sub traffic doesn't need pre-provisioning — the runner adds workers until the backlog drains, then scales back down once it's idle.
gcloud dataflow jobs run word-count-job \
--gcs-location gs://dataflow-templates/latest/Word_Count \
--region us-central1 \
--staging-location gs://my-bucket/staging \
--parameters inputFile=gs://my-bucket/input.txt,output=gs://my-bucket/output \
--num-workers 2 \
--max-workers 20
Two Beam pipelines share the exact same Parse and SumByUser transform code, but one reads from GCS and the other from Pub/Sub. Do you need to rewrite the aggregation logic for the streaming version?
ReadFromText vs ReadFromPubSub), the addition of a windowing stage (streaming data has no natural end, so it needs one), and the write disposition change. You can build and test the logic once against a static batch sample, then point the same pipeline at a live topic.Windowing and Watermarks
An unbounded stream never "finishes," so an aggregation like SumByUser needs a boundary to know when to emit a result — that boundary is a window (fixed, sliding, or session-based, grouping elements by event time, the timestamp the event actually happened, not when it arrived at the pipeline).
The hard part is that event time and processing time diverge — a mobile client goes offline and replays events an hour late, or a network hiccup reorders a batch. Dataflow tracks a watermark: its running estimate of "all data with an event time earlier than T has now arrived." A window only fires (emits its aggregate) once the watermark passes the window's end — but the watermark is an estimate, not a guarantee, so data can still trickle in after a window has already fired. That's late data, and how a pipeline handles it is a real correctness decision, not an edge case to ignore.
FixedWindows(60) window buckets every element by the minute it actually happened, not the minute it was received — so a late-arriving event still lands in the correct historical bucket, not today's.
allowed_lateness configured, the pipeline doesn't discard it — it re-fires that same window with an updated (accumulating) or replacement (discarding) result, and downstream consumers see a correction.
A streaming pipeline uses processing time (wall-clock arrival) instead of event time to bucket a "revenue per minute" aggregation. A batch of events from a flaky mobile client arrives 45 minutes late. Which minute's total do they get counted into?
Dataflow vs. Streaming Straight into BigQuery
bigquery.md already covers the Storage Write API and insertAll — writing rows directly into BigQuery with no separate compute layer in front. That's the simpler, cheaper option whenever it's enough. Add Dataflow in front when the pipeline needs something a single insert call can't express on its own.
bigquery.md) already gives exactly-once delivery and low latency with zero pipeline infrastructure to run or scale. Most simple ingest-and-land use cases stop here.
Why would you add Dataflow in front of BigQuery streaming inserts instead of just streaming directly into BigQuery?
A Representative Pipeline Architecture
graph LR
classDef src fill:#3498db,stroke:#2471a3,color:#fff
classDef compute fill:#e67e22,stroke:#ba6018,color:#fff
classDef sink fill:#27ae60,stroke:#1e8449,color:#fff
classDef state fill:#8e44ad,stroke:#6c3483,color:#fff
PS["Pub/Sub topic<br/>raw click events"]:::src --> DF
subgraph PIPE["Dataflow job — streaming"]
DF["Parse and validate"]:::compute --> WIN["Window: fixed 1-minute,<br/>event time"]:::compute
WIN --> AGG["Stateful aggregation:<br/>per-user session count"]:::state
AGG --> ENRICH["Enrich: side-input lookup<br/>against user profile table"]:::compute
end
ENRICH --> BQ["BigQuery<br/>analytics table"]:::sink
ENRICH --> BT["Bigtable<br/>low-latency lookup by user_id"]:::sink
ENRICH --> GCS["Cloud Storage<br/>raw archive, Avro"]:::sink
This is the shape that justifies Dataflow over a direct insert: one Pub/Sub source feeds a single pipeline that windows, aggregates statefully, enriches against a side input, and then fans out to three different sinks with three different jobs (fast analytical queries, low-latency point lookups, and durable archival) — none of which a single write call could do alone.
Dataproc — Managed Hadoop and Spark
Dataproc runs standard Hadoop/Spark clusters as a managed service — same open-source APIs (Spark, Hive, Pig, MapReduce), no manual node provisioning or cluster software installs. It's the right tool when a Spark or Hadoop codebase already exists, not a reason to start a new pipeline in Spark today.
# Ephemeral pattern: workflow template creates, runs, and tears down automatically
gcloud dataproc workflow-templates create daily-etl --region=us-central1
gcloud dataproc workflow-templates set-managed-cluster daily-etl \
--region=us-central1 \
--cluster-name=daily-etl-cluster \
--num-workers=4 \
--worker-machine-type=n1-standard-4
gcloud dataproc workflow-templates add-job spark \
--workflow-template=daily-etl \
--region=us-central1 \
--step-id=transform \
--class=com.example.SparkTransform \
--jars=gs://my-bucket/jars/transform.jar
gcloud dataproc workflow-templates instantiate daily-etl --region=us-central1
# Cluster is created, the job runs, then the cluster is deleted — one command
# Long-running cluster (interactive / frequent small jobs)
gcloud dataproc clusters create interactive-cluster \
--region=us-central1 \
--num-workers=2 \
--enable-component-gateway \ # web UIs for Spark, Jupyter
--optional-components=JUPYTER
Dataproc vs. rewriting in Dataflow/Beam: reach for Dataproc when there's an existing Spark/Hadoop codebase to migrate, the team already has deep Spark/PySpark expertise, or the job depends on a Spark-specific library (MLlib, GraphX) with no Beam equivalent. Rewriting in Beam is worth the effort for greenfield pipelines, or when you specifically want the batch+streaming unification Beam provides and a fully serverless runner with zero cluster lifecycle to manage.
A team has an existing 2,000-line PySpark job that runs once a night. Is this a good candidate to rewrite in Apache Beam on Dataflow, or to lift into Dataproc as-is?
Orchestration: Cloud Composer vs. Workflows
Neither of these transforms data — they schedule and sequence other things (a Dataflow job, a Dataproc workflow template, a BigQuery load, an HTTP call to some other service). The choice between them comes down to how complex the sequencing is and how much operational overhead is worth paying for that complexity.
# Workflows: a simple HTTP-call-oriented chain (YAML)
main:
steps:
- triggerDataflow:
call: http.post
args:
url: https://dataflow.googleapis.com/v1b3/projects/my-project/locations/us-central1/templates:launch
auth:
type: OAuth2
result: dataflowResult
- waitAndCheck:
call: http.get
args:
url: ${"https://dataflow.googleapis.com/v1b3/projects/my-project/jobs/" + dataflowResult.body.job.id}
auth:
type: OAuth2
result: jobStatus
- notify:
call: http.post
args:
url: https://us-central1-my-project.cloudfunctions.net/notify-slack
body:
status: ${jobStatus.body.currentState}
# Composer: a DAG with heterogeneous, interdependent tasks
from airflow import DAG
from airflow.providers.google.cloud.operators.dataproc import DataprocSubmitJobOperator
from airflow.providers.google.cloud.transfers.gcs_to_bigquery import GCSToBigQueryOperator
from airflow.sensors.filesystem import FileSensor
with DAG("nightly_pipeline", schedule_interval="0 2 * * *") as dag:
wait_for_export = FileSensor(task_id="wait_for_export", filepath="/data/export_ready.flag")
spark_transform = DataprocSubmitJobOperator(task_id="spark_transform", job=SPARK_JOB, region="us-central1")
load_to_bq = GCSToBigQueryOperator(task_id="load_to_bq", bucket="my-bucket", source_objects=["out/*.parquet"],
destination_project_dataset_table="project.dataset.results")
wait_for_export >> spark_transform >> load_to_bq
Cloud Scheduler is the cron layer for either one — it doesn't replace Composer or Workflows, it just fires the trigger. A Scheduler job can hit an HTTP endpoint to start a Workflows execution, or publish to Pub/Sub to kick off a Composer DAG, on a cron schedule (0 2 * * * for "every night at 2am").
# Cloud Scheduler firing a Workflows execution nightly
gcloud scheduler jobs create http nightly-pipeline-trigger \
--schedule="0 2 * * *" \
--uri="https://workflowexecutions.googleapis.com/v1/projects/my-project/locations/us-central1/workflows/nightly-etl/executions" \
--http-method=POST \
--oauth-service-account-email=scheduler-invoker@my-project.iam.gserviceaccount.com
A pipeline needs to: wait for a file sensor, run a Spark job on Dataproc, load results into BigQuery, then trigger three downstream reports with different retry policies — about 15 interdependent tasks in total. Is Workflows a good fit here?
Data Fusion — Visual, No-Code ETL
Data Fusion is a visual pipeline builder (built on the open-source CDAP framework) — drag connectors and transforms onto a canvas, wire them together, and Data Fusion compiles the result down to a Dataproc or Spark job under the hood. No Beam or Spark code to write.
When it's the right call: a data-analyst-heavy team without deep Python/Java/Spark engineers, a pipeline dominated by standard connector-to-connector moves (a database extract, a few field mappings and filters, load into BigQuery) where the built-in connector library covers the sources involved, or when time-to-delivery matters more than hand-tuned performance.
When it's the wrong call: the pipeline needs custom logic that doesn't fit a drag-and-drop transform, the team wants pipeline logic in version control and code review the way hand-written Beam/Spark naturally is, performance or cost needs tight tuning beyond what the visual layer exposes, or the workload is heavy stateful streaming — that's Dataflow's job, not a no-code tool's.
A data analyst needs to pull a weekly CSV extract from an on-prem database, rename a few columns, filter out test accounts, and load it into BigQuery. Is this a good Data Fusion use case, or should it be written as a Beam pipeline?
Which Tool for Which Pipeline Shape
graph TD
classDef direct fill:#7f8c8d,stroke:#616a6b,color:#fff
classDef dataflow fill:#e67e22,stroke:#ba6018,color:#fff
classDef dataproc fill:#9b59b6,stroke:#76448a,color:#fff
classDef fusion fill:#f1c40f,stroke:#b7950b,color:#000
START{"What does the<br/>pipeline actually need to do?"}
START -->|"Light reshape,<br/>one destination table"| DIRECT["Stream directly into BigQuery<br/>(Storage Write API — see bigquery.md)"]:::direct
START -->|"Complex transforms, multiple sinks,<br/>or stateful streaming aggregation"| DF["Dataflow (Apache Beam)"]:::dataflow
START -->|"Existing Spark/Hadoop codebase,<br/>Spark-specific libraries"| DP["Dataproc"]:::dataproc
START -->|"Standard connector-to-connector ETL,<br/>no engineering team needed"| FUS["Data Fusion"]:::fusion
| Criteria | Direct into BigQuery | Dataflow | Dataproc | Data Fusion |
|---|---|---|---|---|
| Transform complexity | None to light (reshape only) | Complex — joins, enrichment, stateful aggregation | Complex, Spark-native (whatever the existing job already does) | Light to medium, connector-driven |
| Team skillset needed | SQL only | Python/Java (Beam SDK) | Spark/Scala/PySpark experience | Low-code — minimal engineering |
| Orchestration complexity | None — fire and forget | Pipeline handles its own windowing; pair with Scheduler/Workflows/Composer for scheduling | Cluster lifecycle to manage; pair with a workflow template or Composer | Built-in scheduling, can also sit under Composer |
| One-off vs. ongoing | Ongoing, continuous | Either — batch or streaming, ongoing production pipelines | Either, but shines for one-off ephemeral migration jobs | Ongoing; less suited to a true one-off |
| Multiple sinks from one source | No — one write target per call | Yes, natively | Possible, but manual to wire up | Yes, via multiple pipeline stages |
| Cost shape | Pay per row written | Pay per worker-hour, autoscaled | Pay per cluster-hour (near-zero if ephemeral) | Pay for the Dataproc/Spark job it compiles to |
A team is deciding between four options for a brand-new pipeline that reads from Pub/Sub, needs a 5-minute rolling deduplication window, and writes to both BigQuery and Bigtable. Which of the four tools in the table actually supports this shape?