UNIT 1 - FOUNDATIONS OF DATA ENGINEERING
1.0 Introduction to Data Engineering
1.1 Definition and Scope of Data Engineering
-
Definition: The practice of designing, building, and maintaining the infrastructure, systems, and processes that enable the collection, storage, processing, and serving of data for analytical and operational use.
-
Scope: Focuses on reliability, scalability, and efficiency of data pipelines. It is the foundational layer that supports Data Science, Business Intelligence (BI), and Machine Learning (ML).
-
Core Objective: Transform raw, often chaotic data into high-quality, accessible, and trusted data assets.
1.2 Role of a Data Engineer
-
Primary Responsibilities:
-
Build and maintain data pipelines (ETL/ELT).
-
Design and implement data storage solutions (Data Lakes, Warehouses).
-
Ensure data quality, reliability, and security.
-
Optimize data systems for performance and cost.
-
Collaborate with Data Scientists, Analysts, and Business Stakeholders.
-
-
Modern Ecosystem Role: Acts as the "plumber" of the data world, ensuring clean, timely data flows to where it's needed, enabling data-driven decision-making.
1.3 Data Engineering Lifecycle
A continuous, cyclical process. The core components are:
-
Data Generation: Creation of data from sources (IoT sensors, applications, logs, user clicks).
-
Data Ingestion: The process of importing/loading data into the storage system. Can be batch (scheduled) or streaming (real-time).
-
Data Storage: Persisting data in systems optimized for purpose (Data Lake for raw, Data Warehouse for structured/processed).
-
Data Processing: Transforming raw data into a usable format. Includes cleaning, enrichment, aggregation (ETL/ELT).
-
Data Serving: Making processed data available and queryable for end-users and applications (via APIs, SQL endpoints, BI tools).
-
Data Consumption: The final use of data by analysts, data scientists, applications, or reports.
DiagramCANVAS: Draw a circular flowchart with 6 labeled stages in order: 1. Generation -> 2. Ingestion -> 3. Storage -> 4. Processing -> 5. Serving -> 6. Consumption. Arrows flow clockwise between them. A feedback loop from Consumption back to Generation/Design indicates iterative improvement.
[!TIP] Exam Focus: Be prepared to draw and explain each stage with a concrete example (e.g., a retail website: clicks -> ingestion -> storage -> processing -> sales dashboard -> business decision).
2.0 Data Architectures and Frameworks
2.1 Types of Data Architectures
| Architecture | Core Idea | Advantages | Disadvantages | Best For |
|---|---|---|---|---|
| Centralized | All data in a single, unified system (e.g., traditional DWH). | Simple governance, single source of truth. | Scalability bottlenecks, vendor lock-in, inflexible. | Small/medium enterprises with stable, structured data. |
| Decentralized | Data resides in multiple, independent systems (data silos). | Domain autonomy, faster iteration per team. | Inconsistent definitions, duplication, integration hell. | Large, complex organizations with distinct business domains. |
| Data Lake | Store all raw data (structured, semi, unstructured) in a single repository (e.g., S3, ADLS). | Schema-on-read, massive scalability, cost-effective, flexible. | Can become a "data swamp" without governance, performance for SQL queries can be poor. | Storing vast amounts of raw data for ML, archival, exploratory analysis. |
| Data Warehouse | Store structured, processed data optimized for SQL analytics (e.g., Snowflake, Redshift). | Schema-on-write, high query performance, strong governance. | Expensive, rigid schema, poor for unstructured data. | Business intelligence, reporting, dashboards on clean, structured data. |
| Lakehouse | Hybrid model combining Data Lake's scalability/flexibility with DWH's ACID transactions & management. | Unified system, supports both BI & ML, cost-effective. | Relatively new, complexity in implementation. | Modern enterprises wanting a single system for all data workloads. |
2.2 Zachman Framework
-
Overview: A six-by-six matrix (framework) for enterprise architecture. It provides a taxonomy to describe an enterprise from different perspectives (Who, What, When, Where, Why, How) and at different abstraction levels (Planner, Owner, Designer, Builder, Subcontractor, Functioning Enterprise).
-
Purpose: To ensure completeness and alignment between business goals and IT systems by answering the six fundamental questions for each stakeholder group.
-
Application: Used as a communication tool and blueprint to organize architectural artifacts, not as a prescriptive methodology.
DiagramCANVAS: Draw a 6x6 grid. Label rows (Perspectives): 1. Executive (Planner), 2. Business Manager (Owner), 3. Architect (Designer), 4. Engineer (Builder), 5. Technician (Subcontractor), 6. System (Enterprise). Label columns (Fundamental Questions): What (Data), How (Function), Where (Network), Who (People), When (Time), Why (Motivation). Cells are empty but represent the intersection of a perspective and a question.
2.3 Lambda Architecture
A unified processing model for handling both batch and real-time data streams.
-
Three Layers:
-
Batch Layer (Speed Layer): Processes the entire historical dataset to create accurate, comprehensive views (e.g., daily). Stores results in a serving layer.
-
Speed Layer (Real-time Layer): Processes new, incoming data in real-time to provide low-latency, approximate views. Compensates for batch layer's latency.
-
Serving Layer: Merges outputs from Batch and Speed layers to serve complete, up-to-date query results to users.
-
-
Process: Raw data is sent to both Batch and Speed layers in parallel. Queries hit the Serving Layer, which combines the best of both.
-
Use Cases: Real-time dashboards, fraud detection, monitoring.
-
Limitations: High complexity (maintaining two separate codebases/logic), operational overhead.
2.4 Kappa Architecture
A simplified alternative to Lambda. It eliminates the batch layer.
-
Core Idea: Process all data as a stream. Historical data is re-processed by replaying the stream from a persistent log (like Kafka).
-
Single Processing Engine: Uses one technology stack (e.g., Apache Flink, Spark Streaming) for both real-time and historical reprocessing.
-
Comparison with Lambda:
| Feature | Lambda Architecture | Kappa Architecture | | :--- | :--- | :--- | | Complexity | High (two pipelines) | Low (one pipeline) | | Code Duplication | Yes (batch & speed logic) | No | | Historical Reprocessing | Native batch layer | Requires stream replay | | Latency | Batch: High, Speed: Low | Consistently Low |
-
When to Use Kappa: When real-time is the primary requirement and the need for complex, heavy batch aggregations is minimal. Preferred for pure streaming applications.
3.0 Data Governance and Maturity
3.1 Data Governance
-
Definition: The overall management of data availability, usability, integrity, and security in an enterprise. It's the people, processes, and policies.
-
Key Principles: Accountability, Standardization, Transparency, Compliance.
-
Core Components:
-
Policies & Standards: Rules for data quality, privacy (GDPR), security, metadata.
-
Roles:
-
Data Owner: Business executive accountable for data asset (defines policies).
-
Data Steward: Subject matter expert who implements and enforces policies (manages quality, definitions).
-
Data Custodian: IT role managing technical infrastructure and security.
-
-
-
Implementation Challenges: Lack of ownership, cultural resistance, tooling complexity, measuring ROI.
3.2 Gartner Data Maturity Model
A five-level model to assess an organization's data management effectiveness.
| Level | Characteristics | Example |
|---|---|---|
| 1. Awareness | Data seen as IT byproduct. No formal strategy. | "We have some spreadsheets." |
| 2. Reactive | Data projects are ad-hoc, siloed, in response to crises. | "The finance team built their own report." |
| 3. Defined | Enterprise-wide strategy exists. Standard tools & processes. | "We have a data governance council and a central DWH." |
| 4. Managed | Data is measured, monitored, and actively managed. Clear metrics. | "We track data quality KPIs and have SLAs." |
| 5. Optimized | Data is a strategic asset, drives innovation, predictive insights. | "Our data products generate revenue and optimize operations in real-time." |
Path to Maturity: Progress is non-linear. Organizations often operate at different levels for different domains. The goal is to move from reactive to proactive, value-driven management.
4.0 Data Lineage and Provenance
4.1 Data Lineage Tracking in DataOps
-
Definition: The lifecycle journey of data—its origins, movements, transformations, and destinations. It answers "Where did this data point come from and how did it get here?"
-
Importance:
-
Compliance & Auditing: (GDPR, HIPAA) Prove data handling.
-
Debugging & Impact Analysis: "Why is this report wrong?" "What breaks if I change this source?"
-
Trust & Transparency: Build confidence in data assets.
-
-
Examples:
-
Financial Transaction: Track from ATM transaction -> core banking system -> nightly batch ETL -> data warehouse -> regulatory report.
-
Healthcare: Track patient lab result from lab instrument -> EHR system -> research data lake -> anonymized ML model.
-
4.2 Pattern-Based Lineage
-
Definition: Automatically inferring lineage by recognizing patterns in code, configuration, or job metadata.
-
How it Works: Tools parse ETL job definitions (SQL scripts, Spark code, Airflow DAGs) to identify source tables, transformations (joins, filters, aggregations), and target tables.
-
Example: An Apache Airflow DAG with a task
transform_salesthat reads fromraw.salesand writes toprod.sales_daily. The tool infers:raw.sales-> [transform_sales] ->prod.sales_daily. -
Pros: Automated, scalable. Cons: Limited to automated jobs; misses manual data edits, spreadsheets.
4.3 Lineage by Data Tagging
-
Definition: Manually or automatically attaching metadata tags (like
PII,Source_System=CRM,Sensitivity=High) to data assets. Lineage is built by connecting tagged assets. -
Approaches:
-
Manual Tagging: Data stewards tag datasets in a catalog (e.g., Alation, Collibra).
-
Automated Tagging: Tools scan data content/statistics to auto-apply tags (e.g., "contains SSN").
-
-
Example: A table
customer_datais taggedSource=CRM,Contains=PII. A downstream tablecustomer_analyticsis taggedDerived_From=customer_data. The lineage link is established via theDerived_Fromtag. -
Pros: Captures non-automated flows (spreadsheets, manual uploads). Cons: Requires discipline, can be incomplete.
4.4 Data Lineage Tools and Techniques
-
Techniques: Automated parsing (SQL, Spark), tag-based correlation, hybrid approaches.
-
Tools: Collibra Lineage, Alation, MANTA, OpenLineage (open-source standard), AWS Glue DataBrew, built-in features in Databricks, Snowflake.
-
Key Technique: Parsing ETL Job Code is the most common automated method for capturing technical lineage.
5.0 Data Storage Patterns
5.1 Data Lake Patterns
-
Design Patterns:
-
Bronze/Silver/Gold Layers (Medallion Architecture):
-
Bronze: Raw, immutable data as ingested. Source of truth for re-processing.
-
Silver: Cleaned, validated, standardized data. Business-level integrity.
-
Gold: Aggregated, business-level data marts/views for specific domains (sales, finance).
-
-
Zone Architecture: Logical separation like
Raw Zone,Enriched Zone,Curated Zone,Serving Zone.
-
-
Merits:
-
Scalability: Store petabytes cost-effectively (object storage).
-
Flexibility: Schema-on-read allows storing anything now, defining structure later.
-
Cost-Effective: Cheap storage tiers (e.g., S3 Glacier).
-
-
Applications: Big Data storage, Machine Learning feature stores, data archival, centralized raw data repository.
5.2 Data Warehouse Patterns
-
Schema Designs:
-
Star Schema: Fact table (central, numeric measures) connected to denormalized dimension tables (descriptive attributes). Simple, fast for queries.
- Example:
sales_fact(sale_id, product_id, store_id, amount) linked todim_product,dim_store.
- Example:
-
Snowflake Schema: Normalized version of Star Schema. Dimensions are broken into multiple related tables. Saves storage but adds join complexity.
-
-
Slowly Changing Dimensions (SCD): Methods to handle changes in dimension table attributes over time.
-
Type 1: Overwrite old value (no history).
-
Type 2: Add new row with version/effective date (full history).
-
Type 3: Add new column for previous value (limited history).
-
5.3 Schema Migration
-
Definition: The process of changing the structure (schema) of a database or data warehouse table (e.g., adding columns, changing data types, splitting tables).
-
Triggers:
-
Business Changes: New product line, regulatory requirement.
-
Technology Upgrades: Migration from on-premise Oracle to cloud Snowflake.
-
Performance Tuning: Normalizing/denormalizing for query speed.
-
-
Example: Migrating a
userstable from on-premise MySQL to cloud Snowflake. Changes:varchar(50)->string, adduser_segmentcolumn, changecreated_attoTIMESTAMP_NTZ. -
Strategies:
-
Big Bang: Switch all applications to new schema at once. High risk, short downtime.
-
Phased (Trickle): Migrate in stages (by table, by user group). Low risk, longer process, requires dual-write or sync mechanisms.
-
6.0 Data Processing and Integration
6.1 Extract, Transform, Load (ETL)
-
Step-by-Step Process:
-
Extract: Pull data from source systems (retail DB:
orders,products,customerstables). -
Transform: Clean, enrich, join. Example: Join
orderswithproductsto getproduct_name, filter out cancelled orders, calculatetotal_amount = quantity * price. -
Load: Write transformed data into the target (Data Warehouse
fact_orderstable).
-
-
Retail Database Example: Nightly batch job extracts new orders from operational PostgreSQL DB, transforms them (adds product category from
productstable, validates email incustomers), and loads into a star-schema DWH for sales reporting. -
Tools: Apache Airflow (orchestration), Informatica PowerCenter, AWS Glue, dbt (transform layer).
6.2 Real-Time Data Processing
-
Concept: Processing data continuously as it arrives (stream) with low latency (seconds/milliseconds).
-
Lambda Architecture for Real-Time (Revisited): Uses the Speed Layer (e.g., Apache Storm, Flink, Spark Streaming) to process the live stream and update a real-time serving layer (e.g., Redis, Cassandra) for low-latency queries. The Batch Layer periodically corrects/improves this view.
6.3 Time-Series Data Pattern Detection (Streaming)
-
Scenario: IoT sensor metrics (temperature, CPU usage), application metrics.
-
Streaming Architecture Components:
-
Ingest: Stream from sensors to message broker (Apache Kafka).
-
Process: Use stream processor (Flink, Spark Streaming) with windowing (tumbling, sliding) to compute aggregates over time.
-
Detect Patterns: Apply algorithms on the windowed stream:
-
Anomaly Detection: Z-score, moving average deviation.
-
Trend Detection: Linear regression on window.
-
Seasonality: Fourier transforms, autocorrelation.
-
-
Serve: Output alerts to dashboard or trigger actions (scale up server).
-
-
Example: Detect server CPU spike:
window(5min).avg(cpu) > 90%-> trigger alert.
6.4 Batch vs. Stream Processing Comparison
| Feature | Batch Processing | Stream Processing |
|---|---|---|
| Data Size | Large, bounded datasets | Continuous, unbounded stream |
| Latency | High (minutes to hours) | Low (ms to seconds) |
| Complexity | Simpler logic | Complex state/time handling |
| Typical Use | Daily reports, historical analysis | Fraud detection, monitoring, real-time dashboards |
| Fault Tolerance | Re-run entire batch | Checkpointing, state recovery |
7.0 Data Operations and Monitoring
7.1 Logging, Monitoring, and Alerting (Observability)
-
Components (Three Pillars):
-
Metrics: Quantitative measurements (CPU %, job duration, row count). Tools: Prometheus, Graphite.
-
Logs: Timestamped records of events/errors. Tools: ELK Stack (Elasticsearch, Logstash, Kibana), Splunk.
-
Traces: End-to-end journey of a request through distributed systems. Tools: Jaeger, Zipkin.
-
-
How it Works: Instrument pipelines to emit metrics/logs. Monitoring tools collect and visualize. Alerting rules (e.g., "Alert if ETL job fails or duration > 2x average") trigger notifications (Slack, PagerDuty).
-
Example: Monitoring an Airflow DAG: Log
task_failure, metricdag_duration_seconds, trace the DAG run ID. Alert ontask_failure.
7.2 Data Quality and Validation
-
Key Checks:
-
Completeness: No nulls in critical fields? Record count matches source?
-
Accuracy: Values within expected range? (e.g.,
agebetween 0-120). -
Consistency:
country_codematchescountry_name? Foreign keys exist? -
Timeliness: Data arrived within SLA?
-
Uniqueness: No duplicate primary keys?
-
-
Implementation: Built into ETL/ELT pipelines (e.g., dbt tests, Great Expectations). Fail the pipeline or quarantine bad data if checks fail.
8.0 Supporting Tools and Protocols
8.1 Webhooks
-
Definition: A user-defined HTTP callback (POST request) triggered by an event in a source system. It's event-driven, push-based.
-
How They Work:
-
User registers a URL (endpoint) with the source system (e.g., GitHub, Stripe).
-
When the event occurs (e.g., code push, payment success), the source system sends an HTTP POST with a payload (JSON/XML) to the registered URL.
-
The receiving endpoint (your service) processes the payload.
-
-
Examples:
-
GitHub: Webhook on
pushevent -> triggers CI/CD pipeline. -
Stripe: Webhook on
payment_intent.succeeded-> updates order status in DB. -
Slack: Incoming webhook -> posts message to channel.
-
-
Advantage: Real-time, efficient (no polling).
8.2 Web Scraping
-
Definition: Automatically extracting data from websites.
-
Techniques:
-
HTML Parsing: Use libraries (
BeautifulSoupin Python,Cheerioin Node.js) to parse DOM. -
API Usage: Prefer official APIs if available (more stable, legal).
-
Headless Browsers: (
Puppeteer,Selenium) for JavaScript-heavy sites. Simulates a real browser.
-
-
Hypothetical Scenario (RGPV Press Releases):
-
Crawl
https://www.rgpvonline.comfor all press release pages. -
On each page, parse HTML to extract
representative_nameandrelease_text. -
Filter releases where
release_textcontains the word "data" (case-insensitive). -
Output list of representative names.
-
-
Ethical & Legal Considerations:
-
Check
robots.txt(guideline, not law). -
Respect
rate limits(don't hammer the server). -
Review Terms of Service (often prohibit scraping).
-
Copyright: Data/ facts may not be copyrighted, but compilation/database might be.
-
Personal Data: Scraping PII may violate privacy laws (GDPR).
-
8.3 Secure Copy Protocol (SCP)
-
Overview: A network protocol based on SSH for secure file transfer between a local and a remote host, or between two remote hosts.
-
How It Works:
-
Uses SSH for authentication and encryption.
-
Connects to remote server via SSH.
-
Starts an
scpprocess on the remote server. -
Transfers files over the encrypted SSH channel.
-
-
Basic Commands:
-
scp file.txt user@remotehost:/path/(local -> remote) -
scp user@remotehost:/path/file.txt .(remote -> local)
-
-
Use Cases in Data Engineering:
-
Securely moving log files or data dumps from application servers to a ingestion server.
-
Initial setup of configuration files on a new cluster node.
-
-
Comparison:
-
vs SFTP: SCP is simpler, faster for single files. SFTP is a full-featured file system protocol (list dirs, resume, permissions). SFTP is generally preferred for robustness.
-
vs Rsync:
rsyncis superior for synchronization (transfers only diffs, resume, checksums).scpalways copies entire file. Use rsync for large datasets or incremental updates.
-
9.0 Case Studies and Applications
9.1 Building a Data Warehouse: University Setup
Step-by-Step Approach:
-
Requirements Gathering: Identify key questions (student performance, course enrollment, faculty workload).
-
Design (Star Schema):
-
Fact Tables:
fact_enrollments(student_id, course_id, faculty_id, grade, credits),fact_financials. -
Dimension Tables:
dim_student(student_id, name, program, admission_date),dim_course,dim_faculty,dim_time.
-
-
ETL Process:
-
Extract: From university's operational systems (Student Info System, HR DB, Finance DB).
-
Transform: Clean data (standardize program names), calculate derived metrics (GPA), slowly changing dimensions for student program changes.
-
Load: Populate dimension tables first (for referential integrity), then fact tables.
-
-
BI & Serving: Connect tools like Power BI/Tableau to the DWH. Build dashboards for Dean (enrollment trends), Faculty (grade distribution), Admin (revenue per program).
9.2 Industry Applications
-
Retail (ETL): Extract from point-of-sale (POS) systems & online logs -> Transform (clean returns, join with product catalog) -> Load into DWH -> Power sales trend dashboards, inventory forecasts.
-
Finance (Real-Time Fraud): Lambda/Kappa Architecture. Stream transaction events -> Speed layer checks against rules (velocity, location) in real-time -> Block suspicious transaction instantly. Batch layer runs nightly complex ML models to refine rules.
-
Healthcare (Data Lake): Ingest massive, diverse data: genomic sequences (FASTA files), EHR records (HL7), medical images (DICOM), research papers (PDF). Store in Data Lake (Bronze). Process to Silver/Gold layers for specific uses: patient 360 views (DWH), ML for drug discovery (feature store).
10.0 Emerging Trends and Challenges
10.1 Shift from ETL to ELT
-
ETL: Transform data before loading into the target (requires powerful staging servers).
-
ELT: Extract and Load raw data into the target (cloud DWH/Lakehouse) first, then Transform using the target's powerful compute.
-
Why? Cloud data platforms (Snowflake, BigQuery) have scalable, separate compute. ELT is simpler, more flexible, leverages target's power, preserves raw data.
10.2 Data Mesh and Domain-Oriented Design
-
Concept: Treat data as a product. Shift from centralized "data platform team" to domain-oriented, decentralized ownership.
-
Principles: Domain-driven data ownership, data as a product, self-serve data infrastructure platform, federated computational governance.
-
Goal: Scale data architecture by aligning data teams with business domains (e.g., "Marketing Data Product", "Supply Chain Data Product").
10.3 Cloud-Native Data Platforms (Snowflake, Databricks)
-
Snowflake: Fully managed, separate storage and compute, multi-cluster shared data, strong SQL DWH. Serverless feel.
-
Databricks: Unified Lakehouse platform built on Apache Spark. Optimized for data engineering, ML, and analytics on data lakes. Uses Delta Lake for ACID transactions.
-
Trend: Move from managing infrastructure (Hadoop on-prem) to consuming managed, scalable cloud services.
10.4 Ethical Considerations in Data Engineering
-
Bias in Pipelines: Can perpetuate or amplify bias if source data is biased (e.g., historical hiring data).
-
Privacy: Designing pipelines that minimize PII collection, enable anonymization/pseudonymization, comply with right-to-be-forgotten.
-
Transparency: Building lineage and documentation to understand how automated decisions are made (explainable AI pipelines).
-
Environmental Impact: Carbon footprint of massive cloud compute and storage. Optimize for efficiency.
-
Data Ownership & Consent: Understanding who owns the data being ingested and processed, especially from third-party sources or user-generated content.
BOXED KEY FORMULAS / DEFINITIONS
-
Data Engineering: \boxed{\text{Design, build, maintain infrastructure for data collection, storage, processing, and serving.}}
-
Data Lineage: \boxed{\text{The lifecycle journey of data: origins, movements, transformations, destinations.}}
-
Lambda Architecture: \boxed{\text{Batch Layer (accuracy) + Speed Layer (latency) + Serving Layer (merged view).}}
-
Kappa Architecture: \boxed{\text{All data processed as a stream; historical data replayed from log.}}
-
ETL Process: \boxed{\text{Extract} \rightarrow \text{Transform} \rightarrow \text{Load}}
-
ELT Process: \boxed{\text{Extract} \rightarrow \text{Load} \rightarrow \text{Transform}}
-
Star Schema: \boxed{\text{Fact Table (measures) + Denormalized Dimension Tables (descriptors).}}
-
SCD Type 2: \boxed{\text{Add new row with version/effective_date to preserve history.}}
-
Webhook: \boxed{\text{Event-triggered HTTP callback (push notification).}}