UNIT 2: DATA ENGINEERING & ARCHITECTURES IN CLOUD COMPUTING
1. FOUNDATIONS OF DATA ENGINEERING
Definition & Scope
Data Engineering is the discipline of designing, building, and maintaining systems and infrastructure for collecting, storing, processing, and serving large volumes of data for analytical and operational use. A Data Engineer ensures data is reliable, accessible, and high-quality for Data Scientists, Analysts, and applications.
Data Engineering Lifecycle (Core Components):
-
Data Generation & Ingestion: Capturing data from sources (APIs, logs, IoT sensors, databases).
-
Data Storage & Processing: Storing data in suitable systems (Data Lakes, Warehouses) and transforming it (batch/stream processing).
-
Data Serving & Consumption: Providing processed data to end-users via dashboards, APIs, or ML models.
-
Data Orchestration & Monitoring: Automating and monitoring pipeline workflows (using tools like Airflow).
DataOps & Modern Practices
DataOps applies DevOps principles (CI/CD, automation, collaboration) to data pipelines to improve speed, quality, and reliability.
Data Lineage Tracking (High-Frequency Topic)
-
Importance: Critical for governance, debugging impact analysis, compliance (GDPR, CCPA), and understanding data flow.
-
Pattern-Based Lineage: Infers lineage by analyzing data access patterns and transformations in pipeline code.
- Example: Analyzing Spark job code to map input
orders_raw→transform_orders()→ outputorders_curated.
- Example: Analyzing Spark job code to map input
-
Lineage by Data Tagging: Uses metadata tags (e.g.,
PII,source_system) attached to datasets/tables. Lineage is built by tracking tag propagation.- Example: Table
customer_datataggedPII. Any downstream table created from it inherits thePIItag automatically.
- Example: Table
-
Automated vs. Manual: Modern systems favor automated lineage (via parsing pipeline code or parsing query logs) over manual documentation for accuracy and scalability.
[!TIP] Exam Focus: Be ready to differentiate Pattern-Based vs. Tagging lineage with clear examples. Lineage is a cornerstone of data governance.
2. DATA ARCHITECTURES & PROCESSING MODELS (Major Exam Focus)
Enterprise Frameworks
- The Zachman Framework: A taxonomy (not a methodology) for enterprise architecture. It uses a 6x6 matrix (What, How, Where, Who, When, Why vs. Planner, Owner, Designer, Builder, Sub-contractor, Functioning Enterprise) to classify architectural artifacts. Provides a holistic view.
Data Warehouse Architectures
-
Traditional vs. Modern Cloud DW:
| Feature | Traditional (On-Prem) | Modern Cloud (Redshift, BigQuery, Snowflake) | | :--- | :--- | :--- | | Scaling | Vertical (bigger servers) | Horizontal (separate compute/storage) | | Cost | High CapEx | Pay-as-you-go OpEx | | Maintenance | In-house, complex | Managed service (auto-patching, backups) | | Flexibility | Rigid schema | Semi-structured support (JSON, VARIANT) |
-
Step-by-Step Data Warehouse Build (University Example):
-
Requirements: Identify needs (student analytics, course enrollment trends).
-
Source Identification: Student Info System (SIS), LMS, Finance DB.
-
Modeling: Design star schema (Fact:
Enrollment_Fact; Dimensions:Student_Dim,Course_Dim,Time_Dim). -
Ingestion: Use batch ETL (e.g., daily) to extract from sources to staging area.
-
Transformation: Clean, integrate, and load into dimensional model in Cloud DW.
-
Serving: Connect BI tools (Tableau, Power BI) to DW for dashboards.
-
Big Data & Lambda-Kappa Architectures
-
Lambda Architecture (Batch Layer, Speed Layer, Serving Layer):
graph LR A[All Data] --> B[Batch Layer<br/>/master dataset]; A --> C[Speed Layer<br/>/real-time stream]; B --> D[Serving Layer<br/>/merged view]; C --> D; D --> E[Queries];-
Batch Layer: Computes accurate, comprehensive views from all historical data (e.g., daily Hadoop/Spark job). Slow but precise.
-
Speed Layer: Computes real-time, incremental views on recent data (e.g., Kafka + Flink/Spark Streaming). Fast but approximate.
-
Serving Layer: Merges batch and speed views to serve queries. Handles ad-hoc queries.
-
Processing Real-Time Data with Lambda:
-
Raw stream → Speed Layer for immediate low-latency alerts (e.g., fraud detection).
-
Same raw stream stored in immutable data store (e.g., HDFS, S3).
-
Batch Layer periodically re-processes all data to correct errors and produce gold-standard views.
-
Serving Layer combines both; queries get latest real-time view corrected by batch.
-
-
Advantages: Fault-tolerant (batch corrects speed layer errors), balanced accuracy/latency.
-
Disadvantages: Complex to build/maintain (two codebases, merging logic).
-
Use Cases: Systems needing both real-time dashboards and accurate historical reports (e.g., e-commerce analytics).
-
-
Kappa Architecture: Simplified, stream-only approach. All processing happens on a single stream.
-
How it works: All data is treated as a stream. To recompute historical views, re-process the entire stream log (stored durably in Kafka, Pulsar) using the same stream processing code.
-
Comparison with Lambda:
| Aspect | Lambda | Kappa | | :--- | :--- | :--- | | Complexity | High (2 layers, 2 codebases) | Low (1 layer, 1 codebase) | | Data Replay | Batch layer on stored data | Replay from immutable log | | Latency | Batch (high) + Speed (low) | Uniformly low (stream-only) | | Use Case Fit | Batch + real-time needs | Pure real-time with historical replay needs |
-
[!TIP] Exam Focus: Draw and explain Lambda's 3 layers. Contrast Lambda vs. Kappa clearly—Kappa's key is the immutable log for replay.
Data Lake & Lakehouse Architectures
-
Data Lake Patterns (Zones):
| Zone | Purpose | Example Data | | :--- | :--- | :--- | | Raw (Bronze) | Immutable landing zone, original format. | JSON logs, CSV dumps, PDFs. | | Curated (Silver) | Cleaned, validated, structured for specific domains. |
customer_cleaned.parquet,orders_enriched.parquet. | | Enterprise (Gold) | Highly refined, business-level aggregates for consumption. |daily_sales_summary,customer_360_view. |-
Merits: Flexibility (schema-on-read), cost-effective storage (cheap object store), scalable for all data types.
-
Applications: Data science exploration, storing raw IoT/clickstream data, archival.
-
-
Data Lake vs. Warehouse vs. Lakehouse:
| Feature | Data Lake | Data Warehouse | Lakehouse | | :--- | :--- | :--- | :--- | | Schema | Schema-on-Read | Schema-on-Write | Schema-on-Write (like DW) | | Data Types | All (raw, unstructured) | Structured only | Structured + Semi-structured | | Workloads | Data science, exploratory | BI, SQL analytics | Unified (BI + DS) | | Cost | Very Low (object store) | High (compute+storage) | Low-Medium (on lake storage) | | ACID | No | Yes | Yes (via Delta/Iceberg/Hudi) |
- Lakehouse (e.g., Delta Lake, Apache Iceberg): Combines flexibility & cost of Data Lake with management & ACID of Data Warehouse on the same storage layer.
Streaming & Time-Series Architectures
-
Processing Time-Series Data & Pattern Detection:
-
Ingest: Stream from sensors/logs → Message Broker (Kafka).
-
Process: Use stateful stream processor (Flink, Spark Streaming).
-
Windowing: Tumble/Hopping/Session windows to group events by time.
-
Pattern Detection: Use Complex Event Processing (CEP) libraries (e.g., Flink CEP) to define patterns like
"temperature > 100 followed by pressure drop within 5 min".
-
-
Store/Serve: Write results to time-series DB (InfluxDB, TimescaleDB) or Data Lake for querying/alerting.
-
-
Core Components:
-
Message Broker: Kafka (durable log, pub/sub).
-
Stream Processor: Flink (true streaming, stateful), Spark Streaming (micro-batch).
-
3. DATA GOVERNANCE, MATURITY & QUALITY
Data Governance
Framework of policies, standards, roles, and processes to ensure data is secure, private, accurate, and usable.
-
Key Roles: Data Owner (business accountable), Data Steward (implements rules), Data Custodian (technical management).
-
Tools: Data Catalog (Alation, Collibra) for discovery; Data Quality tools (Great Expectations); Privacy tools (tokenization, masking).
Gartner Data Maturity Model
Stages of organizational data capability progression:
| Stage | Characteristics | Example |
|---|---|---|
| 1. Awareness | Data seen as IT problem, siloed. | "Reports are manual Excel files." |
| 2. Emerging | First data projects, ad-hoc tools. | "One team uses Tableau for sales." |
| 3. Consolidating | Centralized team, basic governance. | "Established data lake, basic glossary." |
| 4. Defined | Formal processes, standard tools. | "Company-wide data catalog, DQ metrics." |
| 5. Managed | Data treated as product, measurable value. | "Data SLAs, cost-per-query tracking." |
| 6. Optimized | Data-driven culture, AI/ML pervasive. | "Real-time self-service, predictive governance." |
Schema Management
-
Schema Migration: Changing the structure (schema) of a data system (e.g., adding a column, changing data type).
-
Compatibility:
-
Forward Compatibility: Old code can read new data (e.g., adding optional field with default).
-
Backward Compatibility: New code can read old data (e.g., removing optional field).
-
-
Schema Evolution in Streaming (Avro/Protobuf): Use schema registry (Confluent Schema Registry). Writers/readers use schema IDs. Allows safe evolution if compatibility rules (BACKWARD, FORWARD) are enforced.
[!TIP] Exam Focus: Know the 6 Gartner stages with a concrete example for each. For schema migration, always state compatibility requirement.
4. DATA INTEGRATION & INGESTION PATTERNS
ETL/ELT Processes
-
ETL (Extract, Transform, Load):
-
Extract: Pull data from Retail Database (orders, customers, products) to staging area.
-
Transform: Clean (remove nulls), join (orders + customers), aggregate (daily sales), in staging.
-
Load: Write transformed, structured data into Data Warehouse tables (
fact_sales,dim_product).
- Traditional, good for complex transformations on small/medium data.
-
-
Modern ELT (Extract, Load, Transform):
-
Extract & Load: Raw data from sources directly loaded into cloud storage/Data Lake (S3, ADLS) or DW raw tables.
-
Transform: Use cloud DW's powerful compute (BigQuery SQL, Snowpark) to transform data in-place.
- Leverages cloud scalability, faster development, keeps raw data.
-
Ingestion Mechanisms
-
Webhooks (Push-Based, Event-Driven):
-
How they work: Source system HTTP POSTs data to a pre-configured URL (your endpoint) when an event occurs.
-
Use Cases: GitHub repo pushes → trigger CI/CD; Payment gateway → update order status; Slack notifications.
-
Security: Use HTTPS, validate signatures (HMAC), use secret tokens.
-
-
Web Scraping (Pull-Based):
-
Tools:
BeautifulSoup(HTML parsing),Scrapy(full framework),Selenium(dynamic JS). -
Scenario: Scraping rgpvonline.com for "data" press releases:
-
Crawl: Use
Scrapyto follow links to representative pages. -
Parse: For each page, extract
representative_nameandpress_release_text. -
Filter: Check if
"data"(case-insensitive) inpress_release_text. -
Output: Store matching reps' names and URLs in CSV/DB.
-
-
Ethical/Legal: Check
robots.txt, respectrate-limiting, review Terms of Service.
-
-
Batch vs. Streaming Ingestion:
| Batch | Streaming | | :--- | :--- | | Periodic (hourly/daily) | Continuous, per-event | | High throughput, high latency | Low latency, lower throughput per batch | | Simple, fault-tolerant | Complex state management | | ETL classic | Lambda/Kappa speed layer |
5. CLOUD-NATIVE TOOLS & OPERATIONS
Data Transfer & Security Protocols
-
Secure Copy Protocol (SCP):
-
What: Uses SSH for encrypted file transfer between local/remote hosts.
-
vs. SFTP: SCP is simpler, faster for single file copies. SFTP is a full subsystem with more features (resume, directory listing).
-
Command:
scp file.txt user@remote_host:/path/ -
Use Case: Securely moving data files between cloud VM and on-prem server.
-
Observability in Data Pipelines
Logging, Monitoring, and Alerting for pipeline health.
-
Logging: Record structured events (pipeline start/end, row counts, errors). Use JSON format.
- Example:
{"pipeline":"daily_sales","status":"SUCCESS","rows_processed":15000}
- Example:
-
Monitoring: Track key metrics over time.
-
Throughput: Rows/sec, GB/hour.
-
Latency: End-to-end delay, stage duration.
-
Errors: Failure rate, error types.
-
Tools: CloudWatch (AWS), Prometheus (metrics) + Grafana (dashboards), Datadog.
-
-
Alerting: Set thresholds on metrics to notify on failures.
-
Example: "Alert if
pipeline_failure_rate > 5%for 10 mins" or "Alert iflatency > 1 hour". -
Channels: Slack, PagerDuty, Email.
-
[!TIP] Exam Focus: For Logging/Monitoring/Alerting, always state what to measure (metrics) and which tool for which function. SCP is often asked as a short note—know its SSH basis and difference from SFTP.
\boxed{\text{UNIT 2 CORE: Lambda/Kappa, Data Lineage, Data Architectures, Governance Maturity, ETL/ELT, Webhooks/Scraping}}