1. Introduction to Data Engineering
Data Engineering is the discipline of designing, building, and maintaining systems and infrastructure for collecting, storing, processing, and analyzing large volumes of data. Its scope encompasses creating reliable data pipelines that transform raw data into high-quality, usable information for analytics, machine learning, and business intelligence.
The Data Engineering Lifecycle is a continuous, iterative process with five core components:
-
Ingestion: Acquiring raw data from diverse sources (databases, APIs, IoT sensors, logs).
-
Storage: Persisting data in appropriate systems (Data Lakes, Data Warehouses, NoSQL DBs).
-
Processing & Transformation: Cleaning, structuring, and enriching data (ETL/ELT).
-
Serving & Analysis: Making processed data available to consumers (BI tools, data scientists, applications).
-
Monitoring & Maintenance: Ensuring pipeline reliability, performance, and data quality.
[!TIP] Exam Focus: Be prepared to draw and explain this lifecycle diagram. The key is to show it as a cycle, not a linear process, emphasizing continuous feedback and monitoring.
2. Data Architecture Patterns
Types of Data Architectures
| Architecture | Core Idea | Advantages | Disadvantages |
|---|---|---|---|
| Monolithic | Single, integrated system (e.g., traditional RDBMS). | Simple, ACID compliance. | Poor scalability, vendor lock-in, rigid schema. |
| Data Warehouse | Centralized, structured storage for analytics (schema-on-write). | High performance for SQL queries, governed. | Expensive, inflexible for unstructured data, high latency. |
| Data Lake | Store all raw data in native format (schema-on-read). | Scalable, cheap, flexible for all data types. | Can become "data swamp" without governance, query performance issues. |
| Lakehouse | Combines Data Lake scalability with Warehouse structure/management. | Unified storage, ACID transactions, BI/ML support. | Emerging technology, complexity in implementation. |
Lambda Architecture
A hybrid architecture for processing both batch and real-time streams to balance latency, throughput, and fault-tolerance.
-
Batch Layer (Speed Layer): Processes the entire dataset for accurate, comprehensive views (e.g., Hadoop/Spark). Slow but thorough.
-
Speed Layer (Hot Path): Processes recent data streams for low-latency, real-time views (e.g., Kafka/Flink). Fast but approximate.
-
Serving Layer: Merges outputs from Batch and Speed layers to serve unified query results.
[!DIAGRAM: CANVAS]
Title: Lambda Architecture Data Flow
Description: A flowchart with three horizontal layers. At the top, "All Data" flows into the Batch Layer (Master Dataset) and a copy goes to the Speed Layer (Real-time Streams). The Batch Layer outputs a Batch View to the Serving Layer. The Speed Layer outputs a Real-time View to the Serving Layer. The Serving Layer merges these views to answer queries. Arrows show the continuous flow and the periodic recomputation of the batch view.
Kappa Architecture
A simplified alternative to Lambda. All processing is done as a stream. Requires a single, powerful stream processing engine (e.g., Apache Flink, Kafka Streams).
-
Core Idea: Eliminate the separate batch layer. To recompute historical data, simply replay the immutable log of all events through the same stream processing logic.
-
Advantage: Simpler codebase (one processing logic), less operational overhead.
-
Disadvantage: Requires highly scalable stream processors; replaying long histories can be resource-intensive.
[!TIP] Key Difference: Lambda = Batch + Speed (two codebases). Kappa = Only Stream (one codebase, replay log for batch).
Data Lake Patterns & Applications
Patterns:
-
Centralized Data Lake: Single repository for all enterprise data.
-
Distributed Data Lake Mesh: Domain-oriented, decentralized lakes with a federated governance layer.
-
Lakehouse Pattern: Adds a structured, transactional layer on top of the Data Lake (e.g., Delta Lake, Apache Iceberg).
Applications:
-
Machine Learning & Data Science: Store raw, semi-structured data for feature engineering.
-
IoT & Log Analytics: Ingest high-volume, unstructured sensor and application logs.
-
Regulatory Compliance & Archives: Cheap long-term storage of raw data for future audit.
3. Data Governance and Maturity Models
Gartner Data Maturity Model
A 5-stage model assessing an organization's data management sophistication:
-
Unaware: No formal data management. Data is an IT by-product.
-
Emerging: Awareness grows. Siloed projects, basic tools.
-
Bronze: Defined policies, some central governance. Focus on compliance.
-
Silver: Proactive governance. Measurable data quality. Business ownership.
-
Gold: Data is a strategic asset. Embedded governance, advanced analytics, AI-driven insights.
Zachman Framework
A schema (not a methodology) for describing an enterprise's architecture from multiple perspectives. It's a 6x6 matrix:
-
Rows (Perspectives): Planner, Owner, Designer, Builder, Subcontractor, Enterprise.
-
Columns (Questions): What (Data), How (Function), Where (Network), Who (People), When (Time), Why (Motivation).
Each cell defines the artifact (e.g., model, diagram) for that perspective and question. It ensures all aspects of architecture are considered.
Fundamental Principles of Data Governance
-
Treat Data as an Asset: Assign ownership and value.
-
Establish Clear Accountability: Data Owners, Stewards, Custodians.
-
Implement Policy & Standards: For quality, security, privacy (e.g., GDPR).
-
Enable Discovery & Trust: Metadata management, data catalog.
-
Measure & Monitor: Define KPIs for data quality and compliance.
4. Data Lineage and Quality
Data Lineage Tracking in DataOps
Data Lineage is the lifecycle of data: its origins, movements, transformations, and destinations. Critical for impact analysis, debugging, and compliance.
-
Pattern-Based Lineage: Automatically infers lineage by analyzing code/logic patterns in pipelines (e.g., parsing Spark/SQL jobs, Airflow DAGs).
- Example: A tool parses a Spark job's
DataFrametransformations (filter,join,groupBy) to build a dependency graph.
- Example: A tool parses a Spark job's
-
Lineage by Data Tagging: Manually or semi-automatically attaches metadata tags (e.g.,
PII,source_sales_db) to datasets and pipelines. Lineage is derived from tag propagation.- Example: Tagging a raw customer table as
PII. Any downstream dataset created from it automatically inherits thePIItag, showing lineage of sensitive data.
- Example: Tagging a raw customer table as
[!TIP] Best Practice: Use a hybrid approach—automated pattern extraction for technical lineage, supplemented by manual tagging for business context.
Data Quality Assurance and Metrics
DQ Dimensions & Metrics:
| Dimension | Definition | Example Metric |
|---|---|---|
| Completeness | % of expected data present. | (Non-null records / Total records) * 100 |
| Validity | Conformance to defined format/rules. | % of emails matching regex pattern |
| Accuracy | Correctness compared to a real-world source. | Match rate against verified master data. |
| Timeliness | Data availability when needed. | Time lag between source update and warehouse refresh |
| Consistency | Agreement across different copies. | Count of mismatched customer addresses |
| Uniqueness | No unwanted duplication. | Duplicate record rate per entity |
Process: Define rules → Profile data → Monitor metrics → Alert on breaches → Root cause analysis → Remediate.
5. Data Ingestion Techniques
Webhooks
Definition: A user-defined HTTP callback (or small code snippet) that is triggered by a specific event in a source system. It's event-driven, push-based integration.
-
How it works:
-
Client registers a URL (endpoint) with the source application.
-
When the event occurs (e.g., new GitHub commit, Stripe payment), the source app pushes an HTTP POST (with event data) to the client's URL.
-
Client's endpoint receives and processes the payload.
-
-
Example: A GitHub webhook sends a JSON payload to your CI/CD server whenever code is pushed to the
mainbranch, triggering a build.
Web Scraping Methodologies and Tools
Methodology:
-
Identify Target: URL structure, data location (HTML tags, JSON APIs).
-
Fetch: Send HTTP request (GET/POST) to retrieve page source.
-
Parse: Extract structured data from HTML/XML using selectors (CSS, XPath).
-
Store: Save extracted data (CSV, DB, Data Lake).
-
Handle Challenges: Dynamic JS (use headless browsers like Selenium), anti-bot measures (rotate IPs, respect
robots.txt), pagination.
Tools: BeautifulSoup/Scrapy (Python), Puppeteer/Selenium (JS/Python), Apify, ParseHub.
Scenario Solution (RGPV Online):
- Scrape
https://www.rgpvonline.comfor press releases.
- Filter releases containing "data" in title/content.
- Extract representative names associated with those releases.
- Output: List of representatives & links to their "data"-related press releases.
Secure Copy Protocol (SCP)
Definition: A network protocol based on SSH for secure transfer of files between a local and a remote host (or between two remote hosts).
-
How it works: Uses SSH for authentication and encryption. Command:
scp [options] source_file user@remote_host:destination_path. -
Key Feature: Files and passwords are encrypted during transit.
-
Use Case in Data Engineering: Securely moving large log files, backup datasets, or configuration files between on-premise servers and cloud VMs.
-
vs. FTP: SCP is encrypted; standard FTP is not secure (plain text).
SFTP(SSH File Transfer Protocol) is a more feature-rich alternative (resume, directory listing).
6. Data Storage and Processing
Data Warehouse Design: University Setup (Step-by-Step)
-
Requirement Gathering: Identify stakeholders (Admin, Exams, Finance, Research). Key questions: Student performance analysis, fee collection, research grants.
-
Identify Data Sources: Student Info System (SIS), Learning Management System (Moodle), Fee Management, HR, Library.
-
Dimensional Modeling (Star Schema):
-
Fact Tables:
Fact_Student_Results(measures: marks, credits),Fact_Fee_Transactions(amount, date). -
Dimension Tables:
Dim_Student(attributes: roll_no, name, program, batch),Dim_Course,Dim_Time,Dim_Faculty.
-
-
ETL/ELT Design: Extract from source DBs → Clean/transform → Load into DW schema.
-
Technology Selection: Cloud-based (Snowflake, BigQuery, Redshift) for scalability. On-premise (Teradata) if legacy.
-
Build & Test: Develop pipelines, validate data, performance tune.
-
Deploy & Train: Deploy to production, train users (BI tools like Power BI/Tableau).
Extract, Transform, Load (ETL) Process: Retail Example
Source: Operational Retail DB (OLTP) with tables: Orders, Order_Items, Products, Customers.
Target: Data Warehouse for sales analytics.
-
Extract:
-
Full/Incremental load from source tables.
-
Tools:
Sqoop(from RDBMS to HDFS), custom scripts, CDC tools (Debezium).
-
-
Transform:
-
Cleaning: Handle nulls, correct data types.
-
Integration: Join
Orders+Order_Items+Productsto getproduct_name,category. -
Aggregation: Calculate
total_sales_per_product_per_day. -
Derivation: Create
profit_margincolumn.
-
-
Load:
-
Load transformed data into DW fact (
Fact_Sales) and dimension (Dim_Product,Dim_Date) tables. -
Use bulk insert operations, handle slowly changing dimensions (SCD Type 2 for
Dim_Productprice changes).
-
\boxed{\text{ETL Flow: Source Systems} \xrightarrow{\text{Extract}} \text{Staging Area} \xrightarrow{\text{Transform}} \text{Data Warehouse}}
Schema Migration Strategies
Scenario: Evolving Dim_Customer from Version 1 (v1) to v2 (add loyalty_tier column).
-
Backward-Compatible (Additive): Simply add new column
loyalty_tier(nullable). Old applications ignore it. Zero downtime. -
Backward-Incompatible (Destructive): Rename/remove column. Requires coordinated deployment: update all downstream consumers before or during the change. High risk.
-
Parallel Schema (Expand-Contract):
-
Expand: Add new table/column (
Dim_Customer_v2) alongside old. Write to both. -
Migrate: Gradually shift consumers to read from new schema.
-
Contract: Once all consumers migrated, decommission old schema.
-
Real-Time Data Processing with Lambda Architecture
Goal: Process streaming clickstream data for real-time dashboard (last 5 min) + accurate daily reports.
-
Batch Layer (Cold Path):
-
Source: Immutable log (Kafka topic) persisted to HDFS/S3.
-
Process: Daily Spark job computes complete session aggregates, user funnels. Outputs to
Batch_View(e.g., Parquet on S3).
-
-
Speed Layer (Hot Path):
-
Source: Same Kafka stream.
-
Process: Flink/Spark Streaming job computes approximate real-time metrics (e.g., active users last 5 min) using sliding windows. Outputs to
Real_time_View(e.g., Redis, Cassandra).
-
-
Serving Layer:
- Query: Dashboard queries a merge API or a DB that combines
Batch_View(for accuracy) andReal_time_View(for latency). E.g.,final_metric = batch_metric + (realtime_metric - batch_metric_for_same_period).
- Query: Dashboard queries a merge API or a DB that combines
Pattern Detection in Time-Series Streaming Data
Use Case: Detect sudden spike in server CPU usage. Architecture: Kafka (stream) → Flink/Spark Streaming (processing) → Alert/DB.
-
Windowing: Define tumbling/sliding windows (e.g., 1-minute windows).
-
Feature Calculation: Per window, compute metrics:
avg_cpu,std_dev_cpu,max_cpu. -
Pattern Matching:
-
Rule-Based:
IF max_cpu > 95% FOR 3 consecutive windows THEN trigger "Critical Spike". -
ML-Based: Use streaming ML (FlinkML) to compare current window's feature vector against historical baseline (anomaly detection).
-
-
Output: Send alert to Slack/PagerDuty, write anomaly event to
Anomaliestable in DW.
7. Operational Excellence and DataOps
Logging, Monitoring, and Alerting
-
Logging: Capture structured events from pipeline components (Airflow tasks, Spark jobs). Use JSON format, centralize (ELK stack: Elasticsearch, Logstash, Kibana).
-
Monitoring: Track system metrics (CPU, memory, disk I/O) and pipeline metrics (throughput, latency, failure rate, data quality scores). Tools: Prometheus + Grafana, Datadog, CloudWatch.
-
Alerting: Define thresholds on monitored metrics. E.g.,
IF pipeline_failure_rate > 5% for 10 min THEN page on-call engineer. Integrate with PagerDuty/Opsgenie.
[!TIP] Key Metrics to Monitor:
Data Freshness(time since last successful run),Record Count Delta(vs. previous run),Error Rate,Resource Utilization.
DataOps Practices and Cultural Aspects
DataOps applies DevOps principles (CI/CD, automation, collaboration) to data pipelines.
-
Practices:
-
CI/CD for Data: Automated testing (data quality, schema), deployment of pipeline code (dbt, Airflow DAGs).
-
Infrastructure as Code (IaC): Manage cloud resources (S3, Redshift) via Terraform/CloudFormation.
-
Automated Testing: Unit tests for transformation logic, data quality validation tests (Great Expectations).
-
Collaboration: Shared tools (Data Catalog), blameless post-mortems for failures.
-
-
Cultural Aspects: Break silos between Data Engineers, Scientists, and Analysts. Foster ownership, shared responsibility for data quality, and continuous improvement.
Continuous Integration and Deployment for Data Pipelines
-
CI (Continuous Integration):
-
Code commit (SQL, Python, dbt models) triggers pipeline.
-
Run automated tests: linting, unit tests, data quality tests on a dev/staging environment.
-
If tests fail, block merge and notify developer.
-
-
CD (Continuous Deployment):
-
After CI passes, automatically deploy code to production (e.g., update Airflow DAG, dbt project).
-
Use canary deployments or feature flags for risky changes.
-
Post-deployment, monitor production metrics closely.
-
-
Tools: Git (versioning), Jenkins/GitLab CI/GitHub Actions (orchestration), dbt (testing), Terraform (infra).