UNIT 5: Data Engineering
I. Foundations and Lifecycle
Data Engineering: Definition & Scope
-
Definition: The practice of designing, building, and maintaining the systems and infrastructure that enable the collection, storage, processing, and serving of data for analytical and operational use.
-
Evolution: From traditional ETL in data warehouses to modern cloud-native, scalable, and real-time data platforms (Lakehouse, Data Mesh).
-
Lifecycle Core Components:
-
Ingestion: Acquiring data from source systems.
-
Storage: Persisting data in appropriate systems (lakes, warehouses).
-
Processing: Transforming, cleansing, and preparing data.
-
Serving: Making data available to consumers via APIs, dashboards, etc.
-
Consumption: End-users (analysts, data scientists, apps) using the data.
-
Role of Data Engineer
-
Responsibilities: Build and maintain data pipelines, ensure data quality and reliability, optimize performance/cost, implement security/governance.
-
Skills: Programming (Python, SQL, Scala), cloud platforms (AWS/Azure/GCP), big data frameworks (Spark, Hadoop), orchestration (Airflow), data modeling.
-
Impact: Enables data-driven decisions by providing trusted, accessible data.
Five V's of Data
| V | Definition | Example |
|---|---|---|
| Volume | Scale of data (TB, PB, EB) | Social media feeds, IoT sensor data |
| Velocity | Speed of data generation/processing | Real-time stock trades, clickstreams |
| Variety | Different formats & structures | Text, images, JSON, CSV, video |
| Veracity | Data quality, accuracy, trustworthiness | Incomplete records, inconsistent formats |
| Value | Business insights derived from data | Predictive maintenance, customer segmentation |
Types & Sources of Data
-
Structured: Fixed schema, relational DBs (e.g., MySQL tables).
-
Unstructured: No predefined schema (e.g., images, PDFs, videos).
-
Semi-structured: Flexible schema (e.g., JSON, XML, logs).
-
Sources: Internal (CRM, ERP) vs. External (APIs, public datasets, web scraping).
-
Batch vs. Stream: Batch (large volumes, scheduled, e.g., daily sales reports) vs. Stream (continuous, low-latency, e.g., fraud detection).
Data-Driven Decision Making
-
Process: Define objective → Collect relevant data → Process & analyze → Generate insights → Make decision → Monitor outcome.
-
Example: Retailer analyzes real-time sales + inventory + weather data to dynamically adjust pricing and stock levels (demand forecasting).
[!TIP] Exam Focus: Be ready to give specific examples for each "V" and clearly differentiate batch/stream ingestion with use cases.
II. Modern Data Architecture
Key Architectural Paradigms
| Architecture | Core Idea | Pros | Cons |
|---|---|---|---|
| Lakehouse | Combines data lake flexibility with warehouse structure/management | Cost-effective, supports ML & BI, ACID transactions | Relatively new, tooling evolving |
| Data Mesh | Decentralized domain-oriented ownership (data as a product) | Scalable, domain expertise, reduces bottlenecks | Complex governance, cultural shift |
| Data Fabric | Unified layer providing integrated access across disparate sources | Simplifies access, abstracts complexity | Implementation overhead, vendor lock-in risk |
Cloud Platforms (Key Services)
-
AWS: S3 (Storage), Redshift (Warehouse), Glue (ETL), Kinesis (Streaming), EMR (Spark/Hadoop).
-
Azure: Blob Storage, Synapse Analytics, Data Factory, Stream Analytics, Databricks.
-
GCP: Cloud Storage, BigQuery, Dataflow (Apache Beam), Pub/Sub, Dataproc.
Diagram Concept:
DiagramCANVAS: Cloud Data Platform - Show S3/Blob/Storage as central raw zone, with Glue/Dataflow processing into Redshift/BigQuery/Synapse, serving to BI tools and ML services.
Data Lake vs. Data Warehouse
| Feature | Data Lake | Data Warehouse |
|---|---|---|
| Schema | Schema-on-read (flexible) | Schema-on-write (rigid) |
| Data | Raw, all formats (structured/unstructured) | Processed, structured only |
| Users | Data scientists, ML engineers | Business analysts, executives |
| Cost | Low-cost object storage | Higher-cost, optimized compute |
| Use Case | Exploratory analytics, data science | Reporting, BI, dashboards |
- Hybrid Approach: Use lake for raw storage, warehouse for curated business data (Lakehouse pattern).
Scalable Infrastructure Design (Example: E-commerce)
-
Assess: Volume (10M users), Velocity (real-time clicks), Variety (logs, transactions, images).
-
Choose Cloud: AWS for breadth of services.
-
Design:
-
Ingestion: Kinesis (stream), S3 (batch uploads).
-
Storage: S3 (data lake - raw/processed zones), Redshift (curated warehouse).
-
Processing: Glue (ETL), EMR (Spark for large-scale ML).
-
Serving: Redshift Spectrum (query S3), QuickSight (BI).
-
-
Optimize: Auto-scaling, spot instances for EMR, S3 lifecycle policies.
Modern Data Strategies
-
Business Alignment: Start with business outcomes (e.g., "reduce churn by 5%").
-
Agility: Modular pipelines, CI/CD for data, cloud-native services.
-
Cost Optimization: Pay-as-you-go, storage tiering, right-sizing compute, serverless options (Glue, BigQuery).
III. Data Ingestion and Integration
Batch vs. Stream Ingestion
| Aspect | Batch Ingestion | Stream Ingestion |
|---|---|---|
| Data Size | Large volumes | Small, continuous records |
| Latency | High (minutes to days) | Low (milliseconds to seconds) |
| Processing | Scheduled (e.g., nightly) | Continuous |
| Tools | Sqoop, custom scripts, batch Glue jobs | Kafka, Kinesis, Pub/Sub |
| Use Case | Daily sales summary | Real-time fraud detection |
Purpose-Built Ingestion Tools
-
Change Data Capture (CDC): Captures row-level changes in DBs (e.g., Debezium, AWS DMS). Used for near-real-time replication.
-
APIs: REST/gRPC endpoints for structured data pull/push.
-
Messaging Systems: Kafka, RabbitMQ for decoupled, durable event streaming.
Webhooks
-
Concept: User-defined HTTP callbacks triggered by an event (reverse API).
-
Workflow: Event occurs → Source system sends HTTP POST to pre-configured URL (webhook endpoint) → Your service processes payload.
-
Example: GitHub webhook sends event data to your CI/CD server on every
git push.
Web Scraping
-
Techniques: HTML parsing (BeautifulSoup, Scrapy), headless browsers (Selenium) for JS-heavy sites.
-
Scenario (Regulatory Data): Scrape government websites (
rgpvonline.com) for press releases containing "data". Use Pythonrequests+BeautifulSoup, filter by keyword, store in S3.# Pseudo-code response = requests.get(url) soup = BeautifulSoup(response.text, 'html.parser') articles = soup.find_all('div', class_='press-release') for article in articles: if "data" in article.text.lower(): save_to_s3(article) -
Considerations:
robots.txt, rate limiting, legal compliance (terms of service).
CI/CD in Data Pipelines
-
Automation: Git for code, automated testing (data quality checks, schema validation), orchestrated deployment (Airflow, Prefect).
-
Tools: Jenkins, GitLab CI, GitHub Actions, dbt Cloud.
-
Process: Code commit → Run tests (Great Expectations) → Deploy to staging → Validate → Promote to prod.
[!TIP] Common Pitfall: Confusing webhooks (server push) with APIs (client pull). Web scraping must respect
robots.txtand terms of use to avoid legal issues.
IV. Storage Solutions
Data Lake Storage Patterns & Zones
-
Zones:
-
Raw (Bronze): Immutable copy of source data.
-
Processed (Silver): Cleansed, conformed, validated data.
-
Curated (Gold): Business-level aggregates, ready for consumption.
-
-
Merits: Low cost, schema flexibility, supports all data types, central repository.
-
Applications: Data science sandbox, regulatory archive, IoT data hub.
-
Example: Raw JSON logs → Silver Parquet with schema → Gold aggregated daily reports.
Data Warehouse Design: Dimensional Modeling
-
Goal: Optimize for query performance and business understanding.
-
Star Schema: One Fact Table (metrics, foreign keys) connected to multiple Dimension Tables (descriptive attributes).
- Example (Retail):
fact_sales(product_id, store_id, date_id, revenue) linked todim_product,dim_store,dim_date.
- Example (Retail):
-
Snowflake Schema: Normalized dimensions (dimensions have sub-dimensions). Reduces redundancy but adds join complexity.
-
University Setup Example:
-
Fact:
fact_enrollment(student_id, course_id, semester_id, grade). -
Dimensions:
dim_student(name, program, year),dim_course(title, credits, dept),dim_semester(term, year).
-
Cloud Storage Security
-
Encryption:
-
At Rest: Server-side (SSE-S3/KMS), client-side.
-
In Transit: TLS/SSL.
-
-
IAM & Access Controls: Principle of least privilege. Use IAM roles/policies, bucket policies, ACLs (minimize).
-
Secure Copy Protocol (SCP): Command-line tool using SSH for secure file transfer between local & remote host (alternative to insecure FTP).
scp -i mykey.pem file.txt ec2-user@ec2-ip:/path/ -
Diagram Concept:
DiagramCANVAS: Cloud Storage Security Layers - Show data encrypted at rest in S3 bucket, IAM policy attached to user/role, bucket policy allowing only specific VPC endpoint, data in transit via HTTPS/SSL.
V. Data Processing Paradigms
ETL vs. ELT
| Aspect | ETL (Extract-Transform-Load) | ELT (Extract-Load-Transform) |
|---|---|---|
| Sequence | Transform before loading into target | Load raw data first, transform in target |
| Target | Optimized warehouse (schema-on-write) | Data lake/cloud warehouse (schema-on-read) |
| Transformation | Done in staging/ETL server | Pushed to powerful target (e.g., Snowflake, BigQuery) |
| When to Use | Legacy warehouses, complex cleansing needed | Cloud data warehouses/lakes, fast ingestion, iterative exploration |
| Tools | Informatica, Talend, SSIS | dbt, Spark SQL, BigQuery ML |
ETL Process (Retail Example)
-
Extract: Pull
orders,order_items,productsfrom transactional MySQL DB. -
Transform:
-
Clean: Handle nulls, correct data types.
-
Join: Merge tables on keys.
-
Aggregate: Calculate
total_salesper product/day. -
Enrich: Add product category from
dim_product.
-
-
Load: Insert into
fact_daily_salestable in Redshift.
Lambda Architecture (Batch + Speed Layers)
-
Goal: Balance latency & accuracy for big data.
-
Layers:
-
Batch Layer: Processes all historical data (Hadoop/Spark). Creates batch view (accurate, high latency).
-
Speed Layer: Processes real-time stream (Spark Streaming, Flink). Creates real-time view (low latency, approximate).
-
Serving Layer: Merges batch & real-time views to serve queries (e.g., via API).
-
-
Diagram Concept:
DiagramCANVAS: Lambda Architecture - Show incoming data flowing to both Batch Layer (HDFS -> Spark -> Batch View) and Speed Layer (Kafka -> Streaming -> Real-time View). Serving Layer queries combine both views for output.
Kappa Architecture (Simplified Real-Time)
-
Idea: Treat all data as stream. Use a single stream processing engine (e.g., Apache Flink) for both real-time and historical reprocessing.
-
Simplification: Eliminates batch layer. To recompute history, replay stream from持久化 log (Kafka).
-
When: When batch processing logic is a subset of stream logic.
Streaming Analytics Pipeline
-
Components: Source (Kafka) → Stream Processor (Flink/Spark Streaming) → Sink (Data Warehouse/NoSQL/Dashboard).
-
Tools: Kafka (ingestion), Flink/Spark Streaming (processing), Cassandra/Redis (sink), Grafana (visualization).
-
Use Cases: Real-time dashboards, anomaly detection, alerting.
Scaling Considerations for Stream Processing
-
Throughput: Records/sec. Scale by increasing parallelism (partitions, executors).
-
Latency: Time from event to result. Minimize by windowing, state management optimization.
-
Fault Tolerance: Checkpointing (to HDFS/S3), exactly-once semantics, replication (Kafka).
-
State Management: Handle large state ( RocksDB ), savepoints for recovery.
[!TIP] Key Difference: Lambda uses two codebases (batch & speed), Kappa uses one. ELT leverages modern cloud warehouse compute; ETL does work upstream.
VI. Big Data Processing Frameworks
Apache Hadoop
-
HDFS: Distributed file system. Blocks (128MB), replication (3x), NameNode (metadata), DataNodes (storage).
-
MapReduce: Programming model.
map()(filter/sort) →shuffle→reduce()(aggregate). Disk-based, high latency. -
YARN: Resource manager. Separates resource management from scheduling.
-
Architecture Diagram:
DiagramSEARCH: Hadoop HDFS architecture diagram showing NameNode, DataNodes, blocks, and MapReduce flow.
Apache Spark
-
Key Features:
-
In-Memory Processing: Caches data in RAM → 100x faster than MapReduce for iterative algorithms.
-
DAG Scheduler: Creates logical execution plan, optimizes workflow.
-
APIs: RDDs (low-level), DataFrames/Datasets (high-level, optimized).
-
Unified Engine: Batch, Streaming (Spark Streaming), ML (MLlib), Graph (GraphX).
-
-
Advantages over MapReduce:
-
Speed (in-memory vs. disk).
-
Ease of use (high-level APIs in Python/Scala/Java/R).
-
Unified stack for diverse workloads.
-
Better iterative processing (ML, graph).
-
Amazon EMR
-
Managed Service: Runs Hadoop/Spark clusters on AWS.
-
Simplifies: Cluster provisioning, scaling, configuration, monitoring.
-
Integration: Native with S3 (storage), IAM (security), CloudWatch (monitoring).
-
Use Case: Large-scale ETL, log processing, ML on petabyte-scale data without managing infrastructure.
VII. Machine Learning Pipeline Integration
Key Stages of ML
-
Data Collection & Storage: Ingest raw data into lake/warehouse.
-
Data Preprocessing & Feature Engineering: Clean, transform, create features. Most time-consuming, highest impact on performance.
-
Modeling: Select algorithm, train on training set.
-
Evaluation: Test on holdout set, metrics (accuracy, F1, RMSE).
-
Deployment: Serve model as API (real-time) or batch prediction.
-
Monitoring & Retraining: Track drift, performance decay, retrain.
Pre-processing & Feature Engineering
-
Importance: "Garbage in, garbage out." Quality features > complex algorithms.
-
Techniques:
-
Handling missing values (impute, drop).
-
Encoding categoricals (one-hot, label).
-
Scaling/normalization (StandardScaler, MinMax).
-
Creating interaction features, date parts, text vectorization (TF-IDF).
-
-
Impact: Can improve model accuracy by 10-30%+.
AWS SageMaker
-
End-to-End ML Development: Built-in Jupyter notebooks, built-in algorithms, automatic model tuning, one-click deployment.
-
Components:
-
Studio: IDE for ML.
-
Training: Managed training jobs on any algorithm (bring your own or built-in).
-
Processing: Data processing/feature engineering jobs.
-
Deployment: Real-time endpoints (REST API) or batch transform jobs.
-
Monitoring: Detect drift, track metrics.
-
-
Scalability: Automatically scales compute for training/deployment; integrates with S3 for data.
ML Infrastructure on AWS (Diagram Concept)
VIII. Data Governance, Security, and Operations
Data Governance
-
Policies: Rules for data quality, security, privacy (GDPR, HIPAA).
-
Ownership & Stewardship: Data owners (business) define policies; data stewards (technical) implement.
-
Implementation Framework: People + Processes + Technology (catalogs, lineage tools, policy engines).
Data Lineage Tracking
-
Pattern-Based: Automatically infers lineage by analyzing pipeline code (e.g., dbt
manifest.json, Airflow DAGs). -
Lineage by Data Tagging: Manually or automatically tag datasets with source/transformation metadata; build graph from tags.
-
DataOps Context: Critical for impact analysis, debugging, compliance, understanding data flow.
-
Example: If
sales_reporttable is wrong, lineage shows it depends onraw_orders→staging_orders→fact_orders. Pinpoint issue source.
Data Maturity Models
-
Gartner: 5 levels (Awareness → Reactive → Defined → Managed → Optimized). Assess across people, process, technology.
-
Zachman Framework: 6x6 matrix (What, How, Where, Who, When, Why) x (Planner, Owner, Designer, Builder, Subcontractor, Functioning Enterprise). Focuses on enterprise architecture, not just data.
-
Assessment & Roadmap: Evaluate current state, identify gaps, prioritize initiatives (e.g., "implement data catalog first").
Logging, Monitoring, and Alerting
-
Tools: Cloud-native (CloudWatch, Azure Monitor), open-source (Prometheus, Grafana), ELK stack.
-
Metrics: Pipeline success/failure rates, data freshness, record counts, latency, resource utilization (CPU, memory).
-
Incident Response: Alerts (Slack, PagerDuty) → Runbook → Triage → Fix → Post-mortem.
-
Example: Airflow DAG failure → CloudWatch alarm → PagerDuty page → Engineer checks logs, fixes dependency, reruns.
Secure Cloud Storage Practices
-
Encryption: Always enable at-rest (SSE-KMS) and in-transit (HTTPS).
-
Key Management: Use cloud KMS (AWS KMS, Azure Key Vault) for customer-managed keys (CMK). Rotate keys.
-
Access Control: IAM roles over users, bucket policies restricting by VPC/prefix, no public access unless required.
-
Compliance: Enable logging (S3 access logs, CloudTrail), use AWS Config rules for compliance as code.
-
Diagram Concept:
DiagramCANVAS: Secure S3 Bucket - Show bucket with SSE-KMS encryption, bucket policy allowing access only from specific IAM role within a VPC endpoint, CloudTrail logging all API calls to a separate audit bucket.
IX. Advanced Topics and Specialized Concepts
Schema Migration
-
Process: Analyze source/target schemas → Design mapping → Develop transformation logic → Test with subset → Execute (often dual-write) → Validate → Cutover.
-
Challenges: Data type mismatches, null handling, referential integrity, downtime minimization.
-
Tools: Custom scripts, AWS DMS (supports heterogeneous migration), Apache Sqoop.
-
Example: Migrating from on-premise Oracle (VARCHAR2) to cloud PostgreSQL (TEXT). Map data types, handle CLOB/BLOB conversion.
Pattern Detection in Time-Series Streaming Data
-
Techniques: Sliding window aggregations, anomaly detection (statistical: Z-score, ML: Isolation Forest), trend/seasonality detection.
-
Architecture: Stream processor (Flink) maintains stateful windows → Apply detection logic → Output alerts to Kafka topic → Dashboard/Alerting system.
-
Example: Detect spike in server CPU usage (>3σ above 1-hour moving average) within 5-minute window.
Data Wrangling & Data Discovery
-
Data Wrangling (Cleaning/Transformation): Handling missing values, outlier treatment, data type conversion, normalization. Tools: OpenRefine, dbt, Pandas.
-
Data Discovery (Cataloging/Profiling): Automated scanning of data assets to extract schema, statistics (min/max, distinct count), sample values, PII detection. Tools: AWS Glue Data Catalog, Azure Data Catalog, Amundsen.
-
Goal: Make data understandable and trustworthy for analysts.
Real-Time Processing with Lambda Architecture (Step-by-Step)
-
Batch Layer: Hourly Spark job on S3 raw data → computes accurate aggregates → writes to batch view (Redshift).
-
Speed Layer: Flink job consumes Kafka clickstreams → computes recent 5-min counts → writes to real-time view (Redis).
-
Serving Layer: API queries Redis for real-time count + Redshift for historical → merges results → serves to dashboard.
-
Reconciliation: Periodically (daily) batch job corrects any inaccuracies in real-time view.
[!TIP] Exam Focus: Be able to draw and explain Lambda/Kappa architectures. For governance, contrast pattern-based vs. tagging lineage. Schema migration requires a clear, step-by-step plan.