पाठ 17 / 25
Dataflow, Dataproc and 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.
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.
त्वरित जाँच: You have existing Spark jobs to move to Google Cloud with minimal rewriting. Which service fits?
- Dataflow
- Dataproc
- Cloud Tasks
- Firestore
Answer
Dataproc — Dataproc runs managed Spark and Hadoop, so existing Spark code moves with little change.