UNIT 2: Data Engineering
I. Foundations of Data Engineering
Definition and Scope of Data Engineering
Data Engineering is the discipline of designing, building, and maintaining the infrastructure and systems that enable the collection, storage, processing, and serving of data for analytical and operational use cases. It focuses on creating reliable, scalable, and efficient data pipelines.
Scope Includes:
-
Building data platforms (data lakes, warehouses).
-
Developing ETL/ELT pipelines.
-
Ensuring data quality, governance, and security.
-
Optimizing data storage and compute costs.
Data Engineering Lifecycle
A continuous cycle with five core components:
| Component | Purpose | Example |
|---|---|---|
| 1. Data Ingestion | Collecting data from various sources into a staging area. | Ingesting logs from web servers, streaming tweets, or batch CSV dumps from an ERP. |
| 2. Data Storage | Persisting raw and processed data in suitable systems. | Storing raw JSON in a Data Lake (S3) and structured aggregates in a Data Warehouse (Snowflake). |
| 3. Data Processing | Transforming raw data into usable formats (cleaning, enriching, aggregating). | Using Spark to cleanse user records and calculate daily sales aggregates. |
| 4. Data Serving | Making processed data available to consumers via APIs, DBs, or files. | Exposing aggregated metrics through a BI tool (Tableau) or a REST API. |
| 5. Data Consumption | End-users or applications using the served data for decisions/ML. | Data scientists training models on the feature store; analysts building dashboards. |
[!TIP] Exam Focus: Be ready to draw and explain this lifecycle. Link each stage to a concrete tool (e.g., Ingestion: Kafka; Storage: S3/HDFS; Processing: Spark; Serving: Redshift).
II. Data Architectures and Patterns
Types of Data Architectures
| Architecture | Description | Pros | Cons |
|---|---|---|---|
| Batch Processing | Processes large volumes of stored data at scheduled intervals. | Cost-effective, simple, handles huge volumes. | High latency (hours/days), not for real-time needs. |
| Real-Time/Streaming | Processes data continuously as it arrives. | Low latency (seconds), immediate insights. | Complex, higher cost, state management challenges. |
| Hybrid (Lambda/Kappa) | Combines batch and streaming layers to balance latency & accuracy. | Flexible, balances speed & correctness. | Increased complexity, operational overhead. |
Data Lake Patterns
A Data Lake is a vast, raw data repository (often object storage) storing structured, semi-structured, and unstructured data in its native format.
Key Characteristics:
-
Schema-on-Read: Apply structure when reading data, not on write.
-
Scalable & Cheap: Built on object storage (S3, ADLS).
-
ACID Support: Modern lakes (Delta Lake, Iceberg) add transactions.
Design Patterns:
-
Medallion Architecture: Bronze (raw), Silver (cleaned), Gold (business-level) layers.
-
Data Mesh: Domain-oriented decentralized ownership.
Merits & Applications:
-
Merits: Flexibility, cost-effective storage, supports ML/advanced analytics.
-
Applications: Storing IoT sensor data, clickstream logs, unstructured documents for future exploration.
Lambda Architecture
A hybrid architecture for real-time processing with fault tolerance.
Three Layers:
-
Batch Layer (Slow/Accuracy): Processes entire historical dataset (e.g., daily Hadoop/Spark job). Outputs batch views.
-
Speed Layer (Fast/Low-Latency): Processes real-time stream (e.g., Kafka + Flink). Outputs real-time views to compensate for batch layer latency.
-
Serving Layer: Merges batch and real-time views to serve final, up-to-date query results (e.g., via a database that can merge incremental updates).
Step-by-Step Real-Time Processing:
-
Ingest data to both batch (distributed storage) and speed (stream) layers.
-
Batch layer recomputes full dataset every cycle (e.g., nightly).
-
Speed layer computes incremental updates on the fly.
-
Serving layer combines both views on read (query).
[!TIP] Common Pitfall: Lambda is complex to maintain (two codebases, merging logic). Kappa is often preferred now.
Kappa Architecture
A simplified, streaming-only approach. All processing happens on a single, real-time stream.
How it works:
-
All data is treated as a stream.
-
Historical data is replayed from a persistent log (Kafka) to rebuild state.
-
A single stream processing engine (e.g., Flink, Samza) handles both real-time and batch-like reprocessing.
Comparison with Lambda:
| Feature | Lambda | Kappa |
|---|---|---|
| Complexity | High (2 pipelines) | Low (1 pipeline) |
| Codebase | Dual (batch + speed) | Single |
| Latency | Low (speed layer) | Low (pure stream) |
| Reprocessing | Full batch recompute | Replay stream from log |
Streaming Architecture for Time-Series Data & Pattern Detection
Pattern Detection Techniques:
-
Sliding Windows: Aggregate metrics over recent time intervals (e.g., "avg temp last 5 min").
-
Session Windows: Group events by activity gaps (e.g., user session).
-
Anomaly Detection: Use statistical models (Z-score) or ML models on stream (e.g., Flink ML).
-
Complex Event Processing (CEP): Detect sequences of events (e.g., "3 failed logins in 1 min → alert").
Example Architecture:
Sensors → Kafka → Flink (windowing + anomaly detection) → Alert/DB
III. Data Storage Solutions
Data Warehouse Design (University Use Case)
Step-by-Step Approach:
-
Requirement Gathering: Identify key entities (Students, Courses, Faculty, Fees) and KPIs (enrollment, pass rate, revenue).
-
Source Analysis: Map source systems (Student DB, LMS, Finance DB).
-
Schema Design: Choose Star Schema for simplicity & query performance.
-
Fact Tables:
Fact_Enrollment(student_id, course_id, date, grade, fee_paid). -
Dimension Tables:
Dim_Student,Dim_Course,Dim_Date,Dim_Faculty.
-
-
ETL Pipeline Design: Extract from sources, transform (clean, conform dimensions), load into warehouse.
-
Tool Selection: Cloud DWH (Snowflake/BigQuery) or on-prem (Redshift).
Data Lake vs. Data Warehouse
| Feature | Data Lake | Data Warehouse |
|---|---|---|
| Data Type | Raw, all formats (JSON, CSV, logs, images). | Structured, processed, modeled. |
| Schema | Schema-on-Read. | Schema-on-Write. |
| Users | Data scientists, ML engineers. | Business analysts, executives. |
| Cost | Low (object storage). | Higher (compute + storage). |
| Agility | High (flexible). | Low (rigid schema). |
| Use Case | Exploratory analytics, ML, storing raw data. | BI, reporting, dashboards, SQL analytics. |
[!TIP] Selection Rule: Use Lake for raw, diverse, future-unknown data. Use Warehouse for structured, fast SQL on known business metrics. Modern pattern: Lakehouse (Delta/Iceberg on lake storage).
IV. Data Integration and Transformation
ETL Process (Retail Example)
Detailed Workflow:
-
Extract: Pull data from source systems.
-
Sources:
OLTP_MySQL(transactions),CRM_Salesforce(customers),ERP_SAP(inventory). -
Method: Full dump (batch) or CDC (Change Data Capture) via logs.
-
-
Transform: Clean, integrate, and model data.
- Steps: Deduplicate, handle nulls, standardize dates, join customer+transaction, calculate
total_sales = qty * price, aggregate todaily_store_sales.
- Steps: Deduplicate, handle nulls, standardize dates, join customer+transaction, calculate
-
Load: Write transformed data to target.
-
Target:
Retail_DWH(Star Schema:Fact_Sales,Dim_Product,Dim_Store,Dim_Date). -
Method: Full refresh (truncate+insert) or incremental (merge).
-
Tools: Apache Airflow (orchestration), Spark/DBT (transform), Snowflake/BigQuery (target).
Schema Migration
Concept: The process of evolving the structure (schema) of a database or data lake table over time to accommodate new business requirements, without breaking existing pipelines.
Necessity: Business needs change (new attribute, new data source), source systems evolve, regulatory requirements add fields.
Example (Data Lake Schema Evolution):
-
Initial Table:
user_activity(columns:user_id, timestamp, page). -
New Requirement: Track device type.
-
Migration: Add
devicecolumn. Using Delta Lake:ALTER TABLE user_activity ADD COLUMNS (device STRING);-
Old queries (without
device) still work. -
New data must include
devicefield. -
Backfill historical data if needed.
-
[!TIP] Key Point: Schema migration in a lake (schema-on-read) is easier than in a warehouse (schema-on-write). Use table formats (Delta, Iceberg) for safe evolution.
V. Data Governance and Lineage
Data Governance
Principles: Accountability, transparency, security, compliance. Policies & Standards: Data ownership, classification (PII), quality rules (completeness >95%), retention policies. Roles:
-
Data Owner: Business head accountable for data.
-
Data Steward: Defines rules, ensures quality.
-
Data Custodian: IT implements controls (access, backup).
Data Lineage Tracking in DataOps
Tracking data flow from source to consumption, showing transformations and dependencies.
1. Pattern-Based Lineage
-
Mechanism: Parse code/pipeline definitions (SQL, Spark, Airflow DAGs) to infer data flow.
-
Example: An Airflow DAG
extract_sales → transform_sales → load_dwh. Lineage tool parses DAG and shows:Source: sales_db.table → [transform_sales.py: filter, aggregate] → Target: dwh.fact_sales -
Tools: OpenLineage, Marquez.
2. Lineage by Data Tagging
-
Mechanism: Attach metadata tags to datasets/tables at each stage. Lineage is built by tracking tag propagation.
-
Example:
-
Raw table
raw_logstaggedsource:web, PII:no. -
After PII masking job, output
clean_logstaggedPII:masked. -
Lineage shows transformation from
raw_logs→clean_logsviamask_pii_job.
-
-
Tools: Data catalog (Alation, Collibra) with tagging APIs.
Importance:
-
Compliance: Prove data origin and handling (GDPR, CCPA).
-
Debugging: Trace error from report back to faulty source or transformation.
-
Impact Analysis: See what breaks if a source table changes.
VI. Tools and Technologies
Webhooks
Working: A user-defined HTTP callback (POST request) triggered by an event in a source system. It's "server-to-server" push, not polling.
Mechanism:
-
User registers a URL (endpoint) with the source system (e.g., GitHub repo settings).
-
Event occurs (e.g.,
pushto repo). -
Source system sends HTTP POST with JSON payload to the registered URL.
-
Receiving server (your service) processes the payload.
Real-World Examples:
-
CI/CD: GitHub webhook → Jenkins triggers build on code push.
-
Notifications: Stripe webhook → Your app updates order status on
payment.succeeded. -
ChatOps: Jira issue updated → Webhook sends message to Slack channel.
Web Scraping (Scenario: rgpvonline.com)
Techniques & Tools:
-
Libraries:
BeautifulSoup(parse HTML),Scrapy(framework),Selenium(JS-heavy sites). -
Ethical/Legal: Check
robots.txt, respectrate-limiting, use official APIs if available.
Scenario Execution:
-
Identify Target: URL:
https://www.rgpvonline.com/press-releases. -
Scrape: Use
requests+BeautifulSoupto get page, parse for<a>tags containing "data" in text. -
Extract: For each matching link, extract representative name (from link text or nearby element) and release URL.
-
Output: CSV/JSON with
[Representative_Name, Press_Release_URL, Date]. -
Quantify: Count unique representatives. Report: "Found 5 representatives with 12 press releases mentioning 'data'."
Secure Copy Protocol (SCP)
Functionality: Uses SSH for secure file transfer between hosts. Encrypts both authentication and data.
Command Example:
scp -r /local/data/ user@remote-server:/remote/backup/
-
-r: recursive copy. -
Uses SSH keys/password for auth.
Security Features:
-
Encryption (AES, etc.) via SSH.
-
Host key verification prevents MITM.
-
No separate daemon; uses SSH service.
Use Cases in Data Engineering:
-
Securely moving log files from application servers to central ingestion point.
-
Transferring large datasets between on-prem and cloud VPCs (before cloud-native tools).
-
Backup/restore of configuration files.
Logging, Monitoring, and Alerting
Implementation Strategy:
-
Logging: Instrument code to emit structured logs (JSON) to a central system (e.g., ELK stack: Elasticsearch, Logstash, Kibana).
-
Monitoring: Collect metrics (pipeline latency, error rates, resource usage) into a TSDB (Prometheus, CloudWatch).
-
Alerting: Define thresholds (e.g., "alert if pipeline fails > 3 times in 10 min") and notify via Slack/Email/PagerDuty (using Alertmanager, Grafana).
Example for Pipeline Health:
-
Log:
INFO: Loaded 10k rows to DWH in 120s→ Kibana. -
Metric:
pipeline_duration_seconds{stage="transform"}→ Prometheus. -
Alert:
ALERT PipelineFailed IF increase(pipeline_failures_total[15m]) > 2.
VII. Frameworks and Models
Zachman Framework
A six-by-six matrix for enterprise architecture, classifying artifacts by who (stakeholder) and what (aspect).
| What (Data) | How (Function) | Where (Network) | Who (People) | When (Time) | Why (Motivation) | |
|---|---|---|---|---|---|---|
| Planner (Executive) | List of things important to the organization | List of business processes | List of locations | List of organizations | List of events/cycles | List of business goals |
| Owner (Business Manager) | Semantic model (conceptual data model) | Business process model | Business logistics model | Work flow model | Master schedule | Business plan |
| Designer (Architect) | Logical data model | Functional model | Distributed system model | Human interface model | Timing model | Business rule model |
| Builder (Implementer) | Physical data model | Technology process model | Technology logistics model | Presentation architecture | Implementation schedule | Rule design |
| Subcontractor (Technician) | Data architecture | System architecture | Technology architecture | Security architecture | Timing architecture | Rule specification |
| Functioning Enterprise (Actual System) | Actual database | Actual process | Actual network | Actual organization | Actual schedule | Actual strategy |
Application: Provides a comprehensive, structured view to ensure all perspectives are considered in enterprise architecture. Used for gap analysis and communication between business & IT.
Gartner Data Maturity Model
Assesses an organization's maturity in managing data as an asset across 5 levels:
| Level | Name | Characteristics | Example |
|---|---|---|---|
| 1 | Initial | Ad-hoc, siloed, reactive. No formal data strategy. | "We export CSVs from apps when needed." |
| 2 | Developing | Some centralized efforts (data warehouse). Limited governance. | "We have a central BI team and a data warehouse for finance." |
| 3 | Defined | Standardized processes, formal governance, data quality programs. | "We have a data governance council, data catalog, and defined SLAs." |
| 4 | Managed | Measured and controlled. Data is a trusted product. Advanced analytics. | "We track data quality metrics, have data product owners, and use ML for predictions." |
| 5 | Optimized | Data-driven culture, continuous improvement, AI/automation pervasive. | "Data products are monetized, real-time decisions are automated, and data ethics are embedded." |
Assessment & Improvement:
-
Assess: Score each dimension (governance, architecture, quality, etc.) against level criteria.
-
Identify Gaps: e.g., "We are Level 2 in governance (reactive) but Level 3 in architecture."
-
Improve: Target next level. For Level 2→3: Establish a data governance council and data catalog.
[!TIP] Exam Answer Structure: Define model → List all 5 levels with 1-sentence characteristics → Give one concrete example per level → Explain how to use it for improvement (gap analysis).