UNIT 4: DATA ENGINEERING, ARCHITECTURES, AND GOVERNANCE
1. Foundations of Data Engineering
Definition & Scope:
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. Its core responsibility is to provide reliable, scalable, and efficient data pipelines that transform raw data into high-quality, accessible data assets.
The Data Engineering Lifecycle:
A continuous cycle with the following key components:
| Stage | Primary Goal | Key Activities & Technologies |
|---|---|---|
| 1. Data Ingestion | Acquire data from diverse sources. | APIs, databases (CDC), logs, IoT sensors, webhooks. Tools: Apache Kafka, AWS Kinesis, Flume. |
| 2. Data Storage | Persist data durably & cost-effectively. | Choose based on structure & access pattern: Data Lakes (S3, ADLS), Data Warehouses (Redshift, BigQuery), NoSQL (Cassandra). |
| 3. Data Processing | Clean, transform, and enrich raw data. | Batch (Spark, Hive) or Stream (Flink, Spark Streaming). Core: ETL/ELT pipelines. |
| 4. Data Serving | Make processed data available to consumers. | Serve via SQL endpoints, APIs, dashboards (Tableau, Power BI), or ML feature stores. |
| 5. Data Consumption | End-users & applications utilize data. | Business intelligence, reporting, machine learning, operational applications. |
Key Principles & Challenges:
-
Principles: Scalability, Reliability, Idempotency, Observability, Automation.
-
Challenges: Schema evolution, data quality at scale, cost management (especially cloud), security/compliance (GDPR, HIPAA), and managing technical debt in pipelines.
[!TIP] Exam Focus: Be prepared to draw and label the lifecycle diagram. Always connect each stage to a real-world tool or technology (e.g., "Ingestion often uses Kafka for streaming").
2. Data Architectures and Patterns
Lambda Architecture
A hybrid architecture for batch & real-time processing of massive datasets.
-
Batch Layer (Slow Layer): Processes all historical data at rest. Computes accurate, comprehensive views. Uses distributed storage (HDFS, S3) and processing (Hadoop, Spark). Outputs to serving layer.
-
Speed Layer (Fast Layer): Processes real-time data streams to provide low-latency views. Uses stream processing (Flink, Storm). Compensates for batch layer latency.
-
Serving Layer: Merges outputs from batch and speed layers to serve unified, queryable views to users. Typically uses a NoSQL DB (e.g., Cassandra, Druid).
-
Workflow: Batch job runs periodically (e.g., hourly) to recompute full view. Speed layer updates view in real-time. Serving layer merges both.
-
Advantages: Fault-tolerant (batch layer is source of truth), balances latency & accuracy.
-
Disadvantages: High complexity (maintain two codebases), operational overhead, potential for data skew between layers.
Kappa Architecture
A simplified alternative to Lambda. All processing is done as a stream.
-
Core Concept: Use a single, powerful stream processing engine (e.g., Apache Flink) to handle both real-time and historical data. Historical data is replayed from a durable log (Kafka) to rebuild state.
-
Comparison with Lambda:
| Feature | Lambda Architecture | Kappa Architecture | | :--- | :--- | :--- | | Complexity | High (two separate codebases) | Low (single codebase) | | Data Latency | Batch layer high latency | Consistently low latency | | Fault Tolerance | Strong (batch recompute) | Relies on log retention & idempotent processing | | Use Case | Very large, immutable historical data | Real-time analytics where recent data is critical |
-
Advantages: Simpler development & maintenance, unified logic, easier to reason about.
-
Disadvantages: Requires highly capable stream processor; replaying all history for corrections can be expensive.
Data Lake Patterns
A centralized repository storing raw, unprocessed data in its native format (object storage like S3).
-
Characteristics: Schema-on-read, store everything (structured, semi, unstructured), cost-effective, scalable.
-
Common Patterns & "Swamp" Avoidance:
-
Medallion Architecture: A layered approach to improve quality:
Bronze(raw),Silver(cleaned, validated),Gold(business-level aggregates). Key pattern to prevent swamps. -
Data Lakehouse: Combines data lake flexibility with data warehouse management (ACID transactions, BI tools). Implemented via Delta Lake, Apache Iceberg, Hudi.
-
-
Merits: Flexibility, low cost, scalability, supports ML/advanced analytics.
-
Applications: Data science exploration, storing IoT/clickstream logs, archival, as a source for data warehouses.
Types of Data Architectures (Comparative Analysis)
| Architecture | Core Idea | Best For | Key Advantages | Key Disadvantages |
|---|---|---|---|---|
| Data Warehouse | Structured, processed data for SQL BI. | Historical reporting, defined KPIs. | High performance SQL, ACID, mature tooling. | Rigid schema, expensive, poor for raw/unstructured data. |
| Data Lake | Raw data in native format. | Data science, storing everything, cheap archival. | Cheap, flexible, scalable, schema-on-read. | Can become "swamp," poor performance for SQL, no built-in governance. |
| Lakehouse | Data lake + warehouse features. | Unified analytics (BI + ML) on one platform. | ACID, BI performance, cost-effective, open format. | Emerging tech, some complexity in management. |
| Data Mesh | Domain-oriented decentralized ownership. | Large, complex organizations with many domains. | Domain autonomy, product thinking, scalable. | High organizational change, complex governance, new tooling. |
[!TIP] Exam Focus: You will likely be asked to compare at least two architectures (e.g., Lake vs. Warehouse vs. Lakehouse). Use the table structure in your answer.
3. Data Governance, Strategy, and Maturity
Data Governance
The overall management of data availability, usability, integrity, and security.
-
Core Pillars:
-
Policies & Standards: Rules (e.g., "PII must be encrypted") and technical formats.
-
Data Quality: Metrics (completeness, accuracy, timeliness) and monitoring.
-
Data Security & Privacy: Access controls, masking, compliance (GDPR).
-
Data Stewardship: Roles & responsibilities (Data Owners, Stewards) for domain data.
-
Metadata Management: Cataloging, lineage, business glossary.
-
-
Implementation Strategy: Start with a pilot (high-impact domain), define clear RACI matrix, invest in tooling (catalog, quality), and foster a data culture.
Gartner Data Maturity Model
A 5-stage model assessing an organization's data capabilities.
| Stage | Characteristics | Example |
|---|---|---|
| 1. Unaware | No formal data processes; data is an IT by-product. | Reports are manual Excel files from different sources. |
| 2. Emerging | Awareness grows; siloed, reactive efforts begin. | Departmental data marts built without coordination. |
| 3. Consolidating | Centralized team (e.g., BI team) forms; some standards. | Enterprise data warehouse project initiated. |
| 4. Defined | Formal governance, enterprise-wide standards, measured quality. | Data catalog deployed, data quality SLAs monitored. |
| 5. Optimized | Data is a strategic asset; embedded in operations & decisions. | Real-time data products, AI/ML integrated into core workflows. |
Zachman Framework
A schema for enterprise architecture, not a methodology. A 6x6 matrix classifying artifacts.
-
Rows (Perspectives): Planner → Owner → Designer → Builder → Subcontractor → Enterprise (Functioning System).
-
Columns (Interrogatives): What (Data) → How (Function) → Where (Network) → Who (People) → When (Time) → Why (Motivation).
-
Use: Provides a comprehensive, holistic view of an enterprise's architecture, ensuring all aspects are considered from multiple stakeholder viewpoints. Each cell defines a distinct architectural artifact (e.g., for "What" from "Owner" perspective = Semantic Model).
4. Data Movement, Transformation, and Storage
ETL vs. ELT Processes
-
ETL (Extract, Transform, Load): Data is transformed in a staging area before loading into the target warehouse.
-
Example (Retail): Extract sales from POS DB → Clean, aggregate, and join with product data in staging server → Load final fact/dimension tables into Snowflake.
-
Use When: Target system is a traditional warehouse with limited compute; transformations are complex.
-
-
ELT (Extract, Load, Transform): Data is loaded raw into target (often a data lake/lakehouse) and transformed there using its compute.
-
Example: Extract sales logs → Load raw JSON into S3 (Bronze) → Use Spark on the same cluster to create Silver/Gold tables.
-
Use When: Target is a scalable cloud data platform (BigQuery, Snowflake); want to keep raw data; transformations are iterative.
-
Schema Migration
Concept: The process of changing the structure (schema) of a database or data warehouse.
-
Triggers: Business requirement changes, performance issues, tech stack migration (e.g., Oracle to Postgres), normalization/denormalization.
-
Process Steps:
-
Assessment: Analyze current schema, data volume, dependencies.
-
Design: Define new schema, mapping rules, data conversion logic.
-
Implementation: Use tools (AWS DMS, custom scripts) to migrate data and transform to new schema.
-
Validation: Verify data integrity, row counts, application functionality.
-
Cutover & Decommission: Switch applications to new schema, archive old.
-
-
Example: Migrating from a denormalized
orderstable (with customer name repeated) to a normalized schema with separatecustomersandorderstables.
Data Warehouse Implementation (University Case Study)
Step-by-Step Approach:
-
Requirement Gathering: Identify key stakeholders (Admin, Exams, Finance, Research). Define KPIs: student pass rate, faculty workload, research funding.
-
Dimensional Modeling (Kimball): Design Star Schema.
-
Fact Tables:
Fact_Student_Performance(student_id, course_id, time_id, grade, credits). -
Dimension Tables:
Dim_Student(demographics),Dim_Course(dept, level),Dim_Time(semester, year),Dim_Faculty.
-
-
Source System Analysis: Identify sources: Student Info System (SIS), HR DB, Finance DB, Research Portal.
-
ETL/ELT Pipeline Design: Build pipelines to extract from sources, map to dimensions/facts, handle slowly changing dimensions (SCD Type 2 for student program changes).
-
Technology Selection: Cloud DWH (e.g., BigQuery), orchestration (Airflow), BI tool (Looker).
-
Build & Load: Develop pipelines, perform initial historical load.
-
Testing & Validation: Ensure data accuracy, performance testing.
-
Deploy & Train: Deploy to production, train users on BI dashboards.
-
Monitor & Maintain: Set up alerts for pipeline failures, data quality checks.
5. Data Lineage, Provenance, and DataOps
Data Lineage Tracking in DataOps:
Tracks the complete lifecycle of data: origins, movements, transformations, and dependencies.
-
Importance: Debugging pipeline failures, impact analysis (what breaks if schema changes?), compliance/audit (GDPR "right to be forgotten"), understanding data for analytics.
-
High-Level Methodologies: Automated parsing of pipeline code, manual entry, or hybrid approaches.
Pattern-Based Lineage
Infers lineage by analyzing data patterns (value distributions, formats) across datasets without explicit metadata.
-
Technique: Compare statistical signatures (e.g., mean, distinct count, regex patterns) between upstream and downstream datasets to infer transformation logic and data flow.
-
Example: Downstream table
customer_ordershas acustomer_emailcolumn with pattern[a-z0-9][email protected]. Upstreamraw_web_logshas auser_emailcolumn with same pattern. System infers afiltertransformation likely occurred between them. -
Best For: Automated discovery in environments with poor metadata.
Lineage by Data Tagging
Explicitly tags data elements (columns, tables) with provenance metadata at each transformation step.
-
Technique: Each pipeline stage adds tags (e.g.,
source:sis_course_table,transform:join_with_faculty,pipeline:dw_daily_run). Lineage is reconstructed by querying the tag history. -
Example: A column
course_avg_gradein a Gold table is tagged withderived_from: silver.fact_course_enrollments, transformation: avg(grade) group by course_id. -
Best For: Environments with strong governance, where explicit, auditable lineage is required.
Comparison:
| Approach | Accuracy | Effort | Best Use Case |
|---|---|---|---|
| Pattern-Based | Medium (inference can be wrong) | Low (automated) | Large, legacy systems with no existing lineage. |
| Tagging | High (explicit) | High (requires discipline) | Regulated industries (finance, healthcare), new systems. |
6. Real-Time and Streaming Data Processing
Processing Real-Time Data with Lambda Architecture
Implements Lambda's Speed Layer for sub-second latency, while Batch Layer ensures accuracy.
-
Speed Layer Implementation:
-
Ingestion: Use a durable log (Apache Kafka) as the single source of truth for all events.
-
Processing: A stream processor (Apache Flink/Spark Streaming) consumes from Kafka.
-
Performs windowing (e.g., tumbling 1-min windows) and stateful operations (e.g., sessionization, aggregations).
-
Writes incremental, real-time views (e.g., "last 5-min active users") to a low-latency serving store (Redis, Cassandra).
-
-
Handling Late Data: Use watermarks and allowed lateness in windowing to handle out-of-order events.
-
-
Batch Layer (Concurrent): Runs periodic (hourly/daily) jobs over all historical data in the data lake (S3) to recompute accurate, complete views. Outputs to the same serving store (overwriting or merging).
-
Serving Layer Queries: Application queries merge results:
query(real_time_view) UNION ALL query(batch_view, filter=recent)or use a layer that automatically merges (e.g., Druid).
Pattern Detection in Time-Series Data
Identifying anomalies, trends, or recurring sequences in continuous streams.
-
Key Streaming Architecture Components:
-
Windowing: Defines time boundaries for computation (tumbling, sliding, session).
-
Stateful Processing: Maintains state (e.g., last 100 values, running average) across events using keyed streams and state backends (RocksDB).
-
Complex Event Processing (CEP): Libraries (Flink CEP) to define patterns of events (e.g., "temperature > 100°C followed by pressure drop within 10 sec").
-
-
Techniques & Examples:
-
Simple Threshold:
if (value > static_threshold) -> anomaly. Stateless. -
Moving Average / Std Dev: Stateful. Flag if
value > avg_last_100 + 3*stddev. -
Seasonal Decomposition: Stateful. Model trend/seasonality (e.g., using Holt-Winters) and detect residuals beyond confidence interval.
-
Machine Learning: Pre-trained model (Isolation Forest, LSTM) applied in stream for prediction.
-
7. Monitoring, Reliability, and Observability
Logging
Recording discrete events (errors, warnings, info) for debugging.
-
Best Practices:
-
Structured Logging: Use JSON format with fixed fields (
timestamp,level,service,message,trace_id). Essential for parsing. -
Log Levels:
DEBUG(dev),INFO(normal ops),WARN(recoverable),ERROR(failed operation),FATAL. -
Aggregation: Centralize logs using tools (ELK Stack, Splunk, Datadog). Include context (request ID, user ID).
-
-
Role: Primary source for post-mortem analysis of failures.
Monitoring
Tracking system health and performance via metrics.
-
Key Metrics:
-
Throughput: Records/sec, MB/sec.
-
Latency: End-to-end delay, processing time per batch/window.
-
Error Rates: % of failed records, exceptions.
-
Resource Utilization: CPU, memory, disk I/O, network.
-
Backlog: Lag in Kafka consumer groups, queue depth.
-
-
Types:
-
Infrastructure Monitoring: Host/VM/container metrics (Prometheus, CloudWatch).
-
Application Monitoring: Pipeline-specific metrics (custom counters, Spark UI metrics).
-
-
Dashboarding: Create service-specific dashboards (Grafana) showing key metrics over time.
Alerting
Notifying humans when metrics breach thresholds.
-
Setting Thresholds: Base on SLOs (Service Level Objectives). E.g., "99% of batches must complete within 5 mins."
-
Meaningful Alerts: Alert on symptoms (high error rate, backlog growing), not causes. Include context (which pipeline, what job).
-
Avoiding Alert Fatigue:
-
Use multi-level alerts (Warning vs. Critical).
-
Implement alert grouping and suppression (e.g., don't alert every minute on same issue).
-
Escalation Policies: Define who gets called for what severity (PagerDuty, Opsgenie).
-
-
Example: Alert if
pipeline_backlog_bytes > 1GB for 10 mins(Critical) orjob_failure_count > 0(Critical).
[!TIP] The Golden Signal: Focus on Latency, Traffic, Errors, Saturation (4 key signals). Your monitoring should answer: Is it slow? Is it broken? Is it full?
8. Integration Tools, Protocols, and Techniques
Webhooks
An HTTP-based, event-driven mechanism for real-time, push-based communication between applications.
-
How it Works:
-
Subscriber registers a URL (endpoint) with the Provider.
-
When an event occurs (e.g.,
payment.succeededin Stripe), the Provider makes an HTTP POST request to the subscriber's URL with event data (usually JSON). -
Subscriber's endpoint processes the payload.
-
-
Examples:
-
GitHub: POST to your server on
pushevent with commit details. -
Stripe: POST to your server on
invoice.payment_succeededwith payment info. -
Slack: Incoming webhook to post a message to a channel.
-
-
Key Points: Requires a publicly accessible endpoint; must respond quickly (200 OK); must validate signature (security).
Secure Copy Protocol (SCP)
A protocol for secure file transfer between hosts, based on SSH.
-
How it Works: Uses the SSH protocol for authentication and encryption. The
scpcommand invokessshto establish a secure channel and then copies files. -
Command Syntax:
scp [options] source_file user@remote_host:destination_path scp user@remote_host:source_file local_destinationExample:
scp report.csv user@data-server:/data/incoming/ -
Security Advantages over FTP: Encrypts both authentication and data transfer. No separate authentication protocol. Inherits SSH's strong security model (keys, host verification).
Web Scraping
Programmatically extracting data from websites.
-
Methodology:
-
HTTP Request: Send GET/POST to URL (using
requestslibrary). -
Parse Response: Parse HTML/XML content using parser (
BeautifulSoup,lxml). -
Extract Data: Use CSS selectors or XPath to locate and extract specific elements.
-
Store Data: Save to CSV, JSON, database.
-
-
Tools/Libraries:
requests+BeautifulSoup(simple),Scrapy(full framework, async),Selenium(for JS-heavy sites). -
Ethical & Legal Considerations:
-
Check
robots.txt:/robots.txtfile specifies allowed/disallowed paths. -
Review Terms of Service: Many sites prohibit scraping.
-
Be Polite: Add delays (
time.sleep()), identify your scraper inUser-Agent, don't overload servers. -
Copyright: Scraped data may be copyrighted.
-
-
Hypothetical Scenario Application (rgpvonline.com):
-
Goal: Find all representatives with press releases about "data".
-
Steps:
-
Crawl site to find press release pages (look for
/press-releasesor similar). -
For each press release page, extract:
-
Representative name (from page header/bio section).
-
Press release title & content.
-
-
Filter content where text contains keyword "data" (case-insensitive).
-
Output list of
{representative_name, press_release_url, snippet_containing_data}.
-
-
Code Snippet (Conceptual):
for page in press_release_pages: soup = BeautifulSoup(requests.get(page).content) text = soup.get_text().lower() if "data" in text: rep_name = soup.select_one('.representative-name').text print(rep_name, page)
-
[!TIP] Exam Warning: For web scraping questions, always mention ethical considerations (
robots.txt, rate limiting). It's a common marking point.