# Dataflow, Dataproc and Pipelines — Google Cloud Platform

Source: https://www.skillbyai.com/en/gcp/b-pipelines

> Move and transform batch and streaming data into BigQuery.

## Managed engines for data pipelines

**Dataflow** is a serverless runner for **Apache Beam** pipelines: one programming model for **batch and streaming**, with autoscaling workers, windowing, watermarks and exactly-once processing within the pipeline. A classic streaming design is **Pub/Sub → Dataflow → BigQuery**: events arrive on a topic, Dataflow parses, enriches and aggregates them, and results land in tables. Google provides ready-made **Dataflow templates** for common moves. **Dataproc** runs managed **Apache Spark** and Hadoop clusters (or serverless Spark batches), useful when you already have Spark code. For simple ingestion you may need no pipeline at all: Pub/Sub **BigQuery subscriptions** write messages directly to a table, and the **BigQuery Data Transfer Service** loads data from SaaS sources and other clouds on a schedule. **Cloud Composer** is managed Apache Airflow for orchestrating multi-step data workflows.

## A small Beam pipeline: Pub/Sub to BigQuery

Run locally with the DirectRunner for testing, or on Dataflow with `--runner=DataflowRunner`.

```python
import json
import apache_beam as beam
from apache_beam.options.pipeline_options import PipelineOptions

opts = PipelineOptions(streaming=True, project="shop-dev-123456", region="asia-south1")

with beam.Pipeline(options=opts) as p:
    (
        p
        | "Read" >> beam.io.ReadFromPubSub(topic="projects/shop-dev-123456/topics/orders")
        | "Parse" >> beam.Map(lambda b: json.loads(b.decode()))
        | "Keep paid" >> beam.Filter(lambda o: o.get("status") == "paid")
        | "Shape" >> beam.Map(lambda o: {"order_id": o["id"], "amount": o["total"]})
        | "Write" >> beam.io.WriteToBigQuery(
            "shop-dev-123456:shop.paid_orders",
            write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND)
    )
```

## A water treatment plant

Pub/Sub is the inlet pipe, Dataflow the treatment plant that filters and mixes as water flows, and BigQuery the reservoir where everyone draws clean water for analysis.

**Quiz:** You have existing Spark jobs to move to Google Cloud with minimal rewriting. Which service fits?

- [ ] Dataflow
- [x] Dataproc
- [ ] Cloud Tasks
- [ ] Firestore

*Answer:* Dataproc. Dataproc runs managed Spark and Hadoop, so existing Spark code moves with little change.
