GCP Data Pipeline: Pub/Sub, Dataflow, Dataproc, Cloud Composer and data flow for ML
1. GCP Data Pipeline Services
| Service | Type | When to Use |
|---|---|---|
| Pub/Sub | Managed message queue | Event streaming, decouple producers/consumers |
| Dataflow | Managed Apache Beam runner | Unified batch + streaming ETL |
| Dataproc | Managed Spark / Hadoop | Existing Spark/Hadoop workloads, ML at scale |
| Cloud Composer | Managed Apache Airflow | Orchestrate multi-step ML workflows |
| Cloud Storage | Object store | Raw data landing zone, model artifacts |
| BigQuery | Data warehouse | Structured analysis, BigQuery ML |
2. Pub/Sub — Event Streaming
Pub/Sub Architecture:
Data Source → Publisher → [Topic] → Subscription → Subscriber
(IoT devices, (Pull or (Dataflow,
web clicks, Push) Cloud Functions,
logs) BigQuery)
Key concepts:
- Topic: named resource where messages are sent
- Subscription: named resource attached to topic
- Publisher: sends messages to topic
- Subscriber: receives messages from subscription
- At-least-once delivery (not exactly-once by default)
| Feature | Details |
|---|---|
| Message retention | 7 days default (configurable) |
| At-least-once delivery | Idempotent subscribers needed |
| Exactly-once | Available in Pub/Sub Lite (same region) |
| Ordering | Enable message ordering with ordering key |
Exam tip: Pub/Sub → Dataflow → BigQuery is an extremely common pipeline pattern on the exam. Pub/Sub ingests, Dataflow transforms, BigQuery stores and analyzes.
3. Cloud Dataflow — Apache Beam
Dataflow is a managed runner for Apache Beam — a framework for unified batch and streaming processing. No server management required.
| Concept | Description |
|---|---|
| Pipeline | Chain of transform operations |
| PCollection | Distributed data collection (bounded or unbounded) |
| Transform | ParDo, GroupByKey, Combine, Flatten, Partition |
| Windowing | Fixed, Sliding, Session windows for streaming |
| Watermarks | Handle late-arriving data in streaming |
Dataflow Windowing for Streaming ML:
Event stream: ──●──●──●──────●──●──●──────●──●──
Fixed Window (1 min):
├─── [W1] ──┤├─── [W2] ──┤├─── [W3] ──┤
Sliding Window (1 min, slide 30s):
├── [W1] ────┤
├── [W2] ────┤
├── [W3] ────┤
Session Window (2 min gap):
├── [S1] ──────────┤ ├── [S2] ──┤
(user session) (new session)
4. Cloud Dataproc — Managed Spark/Hadoop
| Dataproc Feature | Details |
|---|---|
| Cluster lifecycle | Create in 90 seconds, delete after job — cost efficient |
| Ephemeral clusters | Spin up → run job → shut down (per-job pricing) |
| Preemptible VMs | Use for worker nodes to reduce cost 60-80% |
| Component gateway | Access Jupyter, Zeppelin, Spark UI via browser |
| ML libraries | Spark MLlib, TensorFlow on Spark (TFoS) |
5. Cloud Composer — Workflow Orchestration
Cloud Composer is managed Apache Airflow. Use it to orchestrate multi-step ML pipelines including data ingestion, preprocessing, training, and deployment.
Cloud Composer ML Workflow:
[DAG: daily_ml_pipeline]
Task 1: Extract data from BigQuery
↓
Task 2: Run Dataflow preprocessing job
↓
Task 3: Submit Vertex AI Training Job
↓
Task 4: Evaluate model metrics
↓ (if metrics pass threshold)
Task 5: Deploy to Vertex AI Endpoint
6. Data Pipeline Service Selection
| Scenario | Recommended Service |
|---|---|
| Real-time event streaming ingestion | Pub/Sub |
| Unified batch + streaming ETL (no infra mgmt) | Dataflow (Apache Beam) |
| Migrate existing Spark jobs to GCP | Dataproc |
| Complex ML DAG orchestration | Cloud Composer |
| Stream data into BigQuery | Pub/Sub → Dataflow → BigQuery |
| Serverless data processing (SQL) | BigQuery (ETL via SQL) |
7. Practice Questions
Q1: A company receives millions of IoT sensor events per second from factory equipment. They need to process these events in real time, detect anomalies, and store results in BigQuery. Which pipeline architecture is MOST appropriate?
- A) Dataproc → Spark Streaming → BigQuery
- B) Pub/Sub → Dataflow → BigQuery ✓
- C) Cloud Functions → Cloud SQL
- D) Batch upload to Cloud Storage → BigQuery import
Explanation: Pub/Sub ingests high-volume streaming events reliably. Dataflow processes the stream in real time using Apache Beam (windowing, transformations, anomaly detection). BigQuery stores the results for analysis. This is the canonical GCP streaming analytics pattern.
Q2: A data engineering team has an existing Apache Spark job that processes training data for ML models. They want to migrate it to GCP with minimal code changes. Which service should they use?
- A) Cloud Dataflow
- B) Cloud Dataproc ✓
- C) BigQuery ETL
- D) Cloud Composer
Explanation: Cloud Dataproc supports Apache Spark natively, allowing teams to run existing Spark jobs on GCP with minimal changes. Dataflow uses Apache Beam (different programming model). Dataproc is the lift-and-shift option for Spark workloads.
Q3: A team needs to orchestrate a daily ML pipeline that includes data extraction from BigQuery, preprocessing, Vertex AI training, and deployment if accuracy exceeds 90%. Which service handles this workflow orchestration?
- A) Vertex AI Pipelines
- B) Cloud Dataflow
- C) Cloud Composer ✓
- D) Pub/Sub triggers
Explanation: Cloud Composer (managed Apache Airflow) is designed for complex DAG orchestration across multiple services. It handles scheduling, conditional branching (deploy only if accuracy > 90%), retry logic, and monitoring across heterogeneous services like BigQuery, Dataflow, and Vertex AI.