UNIT 1: DATA ENGINEERING FOUNDATIONS & ARCHITECTURES (Cloud Computing Context)
1.0 Introduction to Data Engineering in the Cloud
Data Engineering is the discipline of designing, building, and maintaining systems for data collection, storage, processing, and serving to enable analytics and machine learning. In cloud-centric organizations, it focuses on leveraging scalable, managed services.
Role of a Data Engineer:
-
Build and maintain data pipelines (ETL/ELT).
-
Ensure data quality, reliability, and security.
-
Optimize cost and performance of cloud data infrastructure.
-
Collaborate with Data Scientists, Analysts, and DevOps.
Key Objectives:
-
Data Availability: Right data, right time.
-
Reliability: Fault-tolerant, recoverable systems.
-
Scalability: Handle volume/velocity growth.
-
Cost-Efficiency: Pay-per-use, auto-scaling.
[!TIP]
Exam Focus: Be ready to define Data Engineering and list the four key objectives. Link each objective to cloud benefits (e.g., scalability via auto-scaling groups).
2.0 The Data Engineering Lifecycle
A cyclic process managing data from origin to consumption.
| Stage | Cloud Services (AWS/Azure/GCP) | Key Considerations |
|---|---|---|
| 2.1 Generation & Acquisition | Kinesis, Pub/Sub, Event Hubs, IoT Core | Sources: Databases, APIs, IoT, logs. Methods: Batch (periodic dumps) vs. Streaming (real-time events). |
| 2.2 Storage | S3, ADLS, GCS (Object Storage) | File Formats: Parquet (columnar, efficient), Avro (row-based, schema evolution), ORC (Hive-optimized). Choices: Data Lakes (raw, flexible) vs. Data Warehouses (structured, optimized). |
| 2.3 Processing | Lambda, Dataflow, Databricks, EMR | Batch: Scheduled, large volumes. Stream: Continuous, low-latency. Serverless: No cluster management (e.g., AWS Lambda for small transforms, Dataflow for streaming). |
| 2.4 Serving & Consumption | Redshift, BigQuery, Snowflake | Data marts (department-specific), BI tools (QuickSight, Power BI), ML platforms (SageMaker, Vertex AI), APIs (REST/gRPC). |
| 2.5 Orchestration & Monitoring | Airflow, Step Functions, Cloud Composer | Workflow DAG scheduling, dependency management. Logging (CloudWatch), Alerting (SNS), Observability (metrics, traces). |
[!TIP]
Common Pitfall: Confusing Data Lake (raw, all data) with Data Warehouse (processed, structured). Lake stores data, warehouse stores information.
3.0 Data Architectures & Patterns
3.1 Lambda Architecture
-
Batch Layer: Processes all historical data (Hadoop/Spark) → accurate, slow.
-
Speed Layer: Processes real-time data (Storm/Flink) → low-latency, approximate.
-
Serving Layer: Merges batch & speed outputs → reconciled view.
-
Cloud Example: S3 (raw) → Spark (batch) + Kinesis + Lambda (speed) → Redshift (serving).
-
Advantages: Fault-tolerant, balances accuracy & latency.
-
Disadvantages: Complex, dual codebases, operational overhead.
3.2 Kappa Architecture
-
Unified Streaming Layer: All data through a single stream (Kafka).
-
Simplification: No batch layer; historical data reprocessed from stream log.
-
Use Case: Real-time analytics where absolute accuracy isn't critical (e.g., live dashboards).
-
Cloud-Native Tools: Kafka + ksqlDB/Flink on GCP/Azure.
3.3 Data Lake Architecture
-
Concept: Centralized repository of raw, unprocessed data in native format.
-
Zones:
-
Raw (Bronze): Immutable source copies.
-
Enriched (Silver): Cleansed, conformed data.
-
Curated (Gold): Business-level aggregates for analytics.
-
-
Cloud Implementation: S3 + Lake Formation (AWS), ADLS + Purview (Azure).
-
Merits: Schema-on-read, cost-effective storage, supports all data types.
-
Challenges: Data Swamp (poor governance), security, performance tuning.
3.4 Data Warehouse Architecture
-
Modern Cloud DWs: Snowflake, Redshift, BigQuery, Synapse.
-
Key Feature: Separation of compute and storage → independent scaling.
-
Advantages: High-performance SQL, built-in optimization, ACID compliance, ideal for structured analytics.
3.5 Comparison of Architectures
| Architecture | Best For | Latency | Complexity | Cloud Fit |
|---|---|---|---|---|
| Lambda | Batch + real-time accuracy | Medium | High | Moderate (multiple services) |
| Kappa | Pure streaming, simplicity | Low | Low-Medium | High (Kafka + Flink) |
| Data Lake | Raw storage, ML, exploration | Varies | Medium | High (object storage + processing) |
| Data Warehouse | Structured analytics, BI | Low | Low | High (managed services) |
[!TIP]
Exam Trick: For "real-time analytics with evolving schema", choose Kappa or Data Lake. For "regulatory reporting with historical accuracy", choose Lambda or Data Warehouse.
4.0 DataOps, Lineage & Governance Foundations
4.1 DataOps Principles
-
CI/CD for Data: Automated testing & deployment of pipelines.
-
Automation: Reduce manual errors (Infrastructure as Code).
-
Collaboration: Data engineers, scientists, analysts.
-
Monitoring: Continuous health checks.
4.2 Data Lineage Tracking
-
Importance: Compliance (GDPR), debugging, impact analysis.
-
4.2.1 Pattern-Based Lineage: Inferred from pipeline code/logic (e.g., Airflow DAGs).
Example: Spark job reads
sales_raw→ transforms → writessales_agg. Lineage auto-generated. -
4.2.2 Lineage by Data Tagging: Tags (PII, sensitivity) attached to datasets; track flow via metadata.
Example: Tag
customersasPII; monitor all downstream tables inheriting tag. -
4.2.3 Tools: AWS Glue (pattern), Azure Purview (tagging), Google Data Catalog.
4.3 Data Governance
-
Policies: Access controls, retention, quality standards.
-
Security: Encryption (at rest/in transit), IAM roles.
-
Quality: Validation rules, profiling (e.g., Great Expectations).
-
Cloud: Use native tools (Lake Formation, Purview) for centralized management.
4.4 Gartner Data Maturity Model
| Stage | Characteristics | Example |
|---|---|---|
| Awareness | Data seen as IT problem | No dedicated data team. |
| Reactive | Ad-hoc requests, siloed tools | Excel-based reports, no documentation. |
| Defined | Standardized processes, data stewards | Centralized data warehouse, basic docs. |
| Managed | Metrics, SLAs, active monitoring | Data quality dashboards, lineage. |
| Optimized | Data as strategic asset, predictive governance | AI-driven data quality, automated compliance. |
4.5 Zachman Framework
-
Enterprise Ontology for architecture descriptions.
-
Six Perspectives (Who): Planner, Owner, Designer, Builder, Subcontractor, Enterprise.
-
Six Aspects (What): What (Data), How (Function), Where (Network), Who (People), When (Time), Why (Motivation).
-
Application: Map cloud services to cells (e.g., "What" + "Builder" = Data model in Snowflake).
[!TIP]
Common Confusion: Lineage (data flow) vs. Governance (policies). Lineage is a tool for governance.
5.0 Implementation Patterns & Techniques
5.1 ETL/ELT Processes
-
Traditional ETL: Transform before loading (on-premise heavy).
-
Cloud-Native ELT: Extract → Load to cloud (S3/DW) → Transform in cloud (scalable compute).
-
Retail Example:
-
Extract: Pull from OLTP (PostgreSQL) via CDC (Debezium).
-
Load: Raw to S3 (landing zone).
-
Transform: Spark/Databricks clean, aggregate → load to Redshift.
-
Serve: Tableau connects to Redshift.
-
5.2 Schema Migration
-
Concept: Changing data structure (e.g., add column, split table).
-
Strategies:
-
Backfill: Update historical data.
-
Dual-Write: Write to old & new schema simultaneously.
-
Expand-Contract: Add new column (expand), deprecate old (contract).
-
-
Example: Add
phone_numbertouserstable.-
Expand: Add nullable column.
-
Dual-write: App writes to both old & new fields.
-
Backfill: Populate
phone_numberfrom legacy logs. -
Contract: Remove old field after validation.
-
5.3 Real-Time & Streaming Processing
-
5.3.1 Time-Series Processing:
-
Windowing: Tumbling (fixed), Sliding (overlap), Session (activity-based).
-
Aggregations: Rolling sums, averages over windows.
Example:
-
$$ \text{5-min tumbling window sum} = \sum_{t}^{t+5\text{m}} \text{sales} $$
-
5.3.2 Pattern Detection:
-
Anomalies: Z-score, moving average deviation.
-
Trends: Linear regression on stream (Apache Flink ML).
Example: Detect IoT sensor spike: if
value > μ + 3σin last 10 min → alert. -
5.4 Building a Data Warehouse: University Case Study
-
Requirements: Track student performance, finance, LMS engagement.
-
Source Systems: Student Info (SQL), LMS (Moodle logs), Finance (SAP).
-
Dimensional Modeling:
-
Fact Tables:
fact_enrollments(student_id, course_id, grade, date). -
Dimension Tables:
dim_students,dim_courses,dim_time.
-
-
Platform Selection: BigQuery (serverless, integrates with G Suite).
-
Phases:
-
Phase 1: Ingest SIS → Bronze layer (S3).
-
Phase 2: Build Silver (cleaned) → Gold (dimensional model).
-
Phase 3: Connect Power BI, train analysts.
-
6.0 Data Acquisition & Integration Tools
6.1 Webhooks
-
Concept: Event-driven HTTP callbacks; source system calls your URL on event.
-
How:
-
Register URL with source (e.g., GitHub).
-
Event occurs (push to repo).
-
Source sends POST JSON to your endpoint.
-
Endpoint triggers pipeline (e.g., start ETL).
-
-
Use Cases: CI/CD triggers, data sync (Stripe → data warehouse), notifications.
6.2 Web Scraping
-
Techniques: HTML parsing (BeautifulSoup), headless browsers (Selenium).
-
Tools: Scrapy (framework), Puppeteer.
-
Ethical/Legal: Check
robots.txt, respectrate-limiting, avoid copyrighted data. -
Example Scenario (RGPV press releases):
# Pseudo-code for scraping rgpvonline.com for page in search_results("site:rgpvonline.com \"data\""): soup = BeautifulSoup(page.content) for rep in soup.find_all("representative"): if "press release" in rep.text: save_to_db(rep.name, rep.text, rep.date)
6.3 Secure Copy Protocol (SCP)
-
Operation:
scp file user@host:/path— copies over SSH. -
Security: SSH encryption, key-based auth.
-
Use Cases:
-
Migrate on-premise logs to cloud VM.
-
Hybrid setups: Transfer files from local server to AWS EC2.
-
-
Cloud Note: Prefer AWS DataSync or Azure File Sync for large-scale, managed transfers.
7.0 Operational Excellence: Logging, Monitoring & Alerting
7.1 Importance
-
Reliability: Detect failures early (e.g., pipeline stall).
-
Cost Control: Spot runaway resources.
-
Debugging: Trace data errors to source.
7.2 Logging
-
Centralized Aggregation: CloudWatch (AWS), Stackdriver (GCP), Azure Monitor.
-
Structured Logging: JSON format with
timestamp,level,pipeline_id,message. -
Log Levels: DEBUG, INFO, WARN, ERROR, FATAL.
7.3 Monitoring
-
Key Metrics:
-
Latency: Time per pipeline stage.
-
Throughput: Records/sec.
-
Error Rate: Failed tasks/total.
-
Data Volume: Bytes in/out.
-
-
Dashboarding: Grafana (custom), CloudWatch Dashboards.
7.4 Alerting
-
Thresholds: e.g., "Alert if error rate > 5% for 5 min".
-
Channels: SNS (AWS), PagerDuty, Slack.
-
Prevent Fatigue: Group alerts, set severity levels, avoid false positives.
7.5 Integrated Observability
-
AWS Glue: Monitor job runs, DPU usage in Glue Console.
-
Databricks: Cluster metrics, job spark UI.
-
Cloud-Native: Services emit metrics to CloudWatch/Monitor automatically.
[!TIP]
Exam Example: "How would you monitor a streaming pipeline?" → Answer: Track lag (Kafka consumer delay), throughput, serialization errors; dashboard in CloudWatch; alert on lag > threshold.
8.0 Data Lake Patterns & Advanced Considerations
8.1 Medallion Architecture (Enhanced Data Lake)
-
Bronze: Raw, immutable ingestion layer (S3 raw zone).
-
Silver: Validated, cleansed, conformed data (Parquet, partitioned).
-
Gold: Business-level aggregates, ready for consumption (dim/fact tables).
-
Cloud: Implement via Delta Lake (Databricks) or Apache Iceberg.
8.2 Merits
-
Schema-on-Read: Flexibility for new data sources.
-
Cost: Cheap object storage (S3 $0.023/GB).
-
Diverse Data: Store logs, JSON, images, CSVs together.
8.3 Applications
-
Data Science Sandbox: Raw data for exploration.
-
Raw Archive: Compliance storage (7+ years).
-
IoT Ingestion: High-volume sensor data.
8.4 Challenges & Solutions
| Challenge | Solution |
|---|---|
| Data Swamp (unorganized) | Enforce medallion zones, metadata catalog (Glue). |
| Security | Fine-grained IAM (S3 bucket policies), encryption (KMS). |
| Metadata Management | Use Apache Hive Metastore or cloud catalog. |
| Performance | Partitioning (by date), Z-ordering (Delta Lake), file sizing (128 MB–1 GB). |
[!TIP]
Key Distinction: Data Lake = raw storage + processing. Data Lakehouse = Lake + Warehouse features (transactions, BI). Mention Delta Lake as example.