Chuyển đến nội dung chính

Lesson 3: Data Pipeline — Dataflow, Pub/Sub, Dataproc

Apache Beam on Dataflow for batch/streaming ETL. Pub/Sub for event-driven pipelines. Dataproc for Spark. Cloud Composer (Airflow) for orchestration.

GCP Data Pipeline Architecture

GCP Data Pipeline: Pub/Sub, Dataflow, Dataproc, Cloud Composer and data flow for ML

1. GCP Data Pipeline Services

ServiceTypeWhen to Use
Pub/SubManaged message queueEvent streaming, decouple producers/consumers
DataflowManaged Apache Beam runnerUnified batch + streaming ETL
DataprocManaged Spark / HadoopExisting Spark/Hadoop workloads, ML at scale
Cloud ComposerManaged Apache AirflowOrchestrate multi-step ML workflows
Cloud StorageObject storeRaw data landing zone, model artifacts
BigQueryData warehouseStructured 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)
FeatureDetails
Message retention7 days default (configurable)
At-least-once deliveryIdempotent subscribers needed
Exactly-onceAvailable in Pub/Sub Lite (same region)
OrderingEnable 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.

ConceptDescription
PipelineChain of transform operations
PCollectionDistributed data collection (bounded or unbounded)
TransformParDo, GroupByKey, Combine, Flatten, Partition
WindowingFixed, Sliding, Session windows for streaming
WatermarksHandle 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 FeatureDetails
Cluster lifecycleCreate in 90 seconds, delete after job — cost efficient
Ephemeral clustersSpin up → run job → shut down (per-job pricing)
Preemptible VMsUse for worker nodes to reduce cost 60-80%
Component gatewayAccess Jupyter, Zeppelin, Spark UI via browser
ML librariesSpark 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

ScenarioRecommended Service
Real-time event streaming ingestionPub/Sub
Unified batch + streaming ETL (no infra mgmt)Dataflow (Apache Beam)
Migrate existing Spark jobs to GCPDataproc
Complex ML DAG orchestrationCloud Composer
Stream data into BigQueryPub/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.