UNIT 5: DATA ENGINEERING ARCHITECTURES, PROCESSES, AND GOVERNANCE
I. FOUNDATIONS OF DATA ENGINEERING
Data Engineering Lifecycle
Data Engineering is the discipline of designing, building, and maintaining systems for collecting, storing, processing, and serving data for analytical and operational use.
The lifecycle consists of five core components:
-
Ingestion: Acquiring data from source systems (e.g., APIs, logs, IoT sensors) into the data platform.
-
Storage: Persisting data in appropriate systems (e.g., Data Lakes, Data Warehouses, NoSQL databases).
-
Processing: Transforming raw data into usable formats (e.g., cleansing, aggregation, enrichment) via batch or streaming jobs.
-
Serving: Making processed data available to consumers via query engines, APIs, or BI tools.
-
Consumption: End-users (analysts, data scientists, applications) accessing data for insights and decisions.
[!TIP] Exam Focus: Be prepared to draw this lifecycle diagram and explain each stage with a real-world example (e.g., "social media app -> ingestion -> storage -> processing -> serving -> analyst dashboard").
Gartner Data Maturity Model
A framework to assess an organization's data management sophistication across five levels:
| Level | Description | Organizational Example |
|---|---|---|
| 1. Unaware | No formal data strategy; data is an afterthought. | Data stored in individual spreadsheets; no central reporting. |
| 2. Opportunistic | Ad-hoc projects; siloed solutions; inconsistent results. | A marketing team builds a standalone campaign analysis tool. |
| 3. Systematic | Defined processes, roles, and technologies; some integration. | Enterprise Data Warehouse (EDW) exists with standard ETL pipelines. |
| 4. Differentiating | Data is a strategic asset; advanced analytics (ML/AI); strong governance. | Real-time personalization engine driving core revenue. |
| 5. Transformational | Data-driven business model; pervasive data culture; automated insights. | Company's primary product is data-driven (e.g., Netflix recommendations). |
Data Governance
The overarching framework of policies, standards, processes, and roles to ensure data is available, usable, consistent, trusted, and secure.
-
Key Roles:
-
Data Owner: Business executive accountable for data assets (defines policies).
-
Data Steward: Subject matter expert who implements policies (manages quality, definitions).
-
Data Custodian: IT role responsible for technical security and storage.
-
-
Importance:
-
Compliance: Meets regulatory requirements (GDPR, HIPAA).
-
Quality: Improves decision-making reliability.
-
Security: Enforces access controls and protects sensitive data.
-
Value: Enables discovery and trusted use of data assets.
-
II. DATA ARCHITECTURES & PATTERNS
Types of Data Architectures
| Architecture Type | Advantages | Disadvantages | Typical Use Case |
|---|---|---|---|
| Centralized | Single source of truth; easier governance & security. | Single point of failure; scalability bottlenecks; less flexible. | Traditional Enterprise Data Warehouse (EDW). |
| Decentralized (Federated) | Domain autonomy; faster innovation; scalable. | Data silos; inconsistent definitions; complex integration. | Microservices with independent databases. |
| Monolithic | Simple to develop/deploy initially; strong consistency. | Inflexible; hard to scale/update; technology lock-in. | Legacy on-premise ERP systems. |
| Microservices-based | Independent scaling; technology diversity; resilience. | Network latency; complex distributed transactions; operational overhead. | Cloud-native applications (e.g., e-commerce platform). |
| Cloud-Native | Elastic scalability; pay-as-you-go; managed services. | Vendor lock-in; egress costs; security configuration complexity. | Modern data platforms on AWS/Azure/GCP. |
| On-Premise | Full control; predictable costs (CapEx); data sovereignty. | High upfront cost; limited scalability; maintenance burden. | Highly regulated industries (defense, banking). |
Data Lake Patterns
A Data Lake is a vast, raw data repository storing structured, semi-structured, and unstructured data in its native format.
-
Design Principles:
-
Store Everything: Capture all raw data, regardless of immediate use.
-
Schema-on-Read: Apply structure when data is read/queried, not on write.
-
Immutable Storage: Raw data is never altered; transformations create new copies.
-
Separation of Storage & Compute: Scale resources independently.
-
-
Merits:
-
Scalability: Store petabytes cost-effectively (e.g., on object storage like S3).
-
Flexibility: Adapt to new data types/analytics without re-ingestion.
-
Cost-Effective: Cheap storage tiering (hot, cool, archive).
-
-
Applications: Advanced analytics, machine learning model training, data archival, exploratory data analysis.
Streaming & Hybrid Architectures
Lambda Architecture
A hybrid model for both batch and real-time processing of large datasets.
| Layer | Purpose | Technology Examples | Process |
|---|---|---|---|
| Batch Layer | Computes accurate results from all historical data. | Hadoop, Spark | Stores master dataset (immutable). Runs long-running jobs to generate batch views. |
| Speed Layer | Computes real-time views on recent data to compensate for batch latency. | Apache Flink, Storm, Kafka Streams | Processes data streams with low latency. Updates real-time views. |
| Serving Layer | Serves merged batch & real-time views to applications. | Cassandra, HBase, Druid | Indexes batch views; merges with real-time views for query response. |
[!TIP] Common Pitfall: Lambda is complex to maintain (two separate codebases for batch & speed layers).
Kappa Architecture
A simplified, streaming-only simplification of Lambda.
-
Core Idea: Treat all data as a stream. Use a single, powerful stream processing engine for both real-time and historical (re-playable) data.
-
How it works:
-
All incoming data is written to a durable, distributed log (e.g., Apache Kafka).
-
A single stream processing system (e.g., Apache Flink) reads from this log.
-
For historical analysis, the log is re-processed from the beginning.
-
Processed results are served from a scalable serving layer (e.g., Elasticsearch).
-
-
Comparison with Lambda:
-
Kappa: Simpler (one codebase), but requires stream processing to handle both low-latency and large-scale batch workloads.
-
Lambda: More complex (two codebases), but allows using optimal tools for batch (accuracy) and speed (latency).
-
Enterprise Architecture Frameworks: Zachman Framework
A schema for organizing enterprise architecture artifacts into a 6x6 matrix.
-
Rows (Perspectives): Planner, Owner, Designer, Builder, Subcontractor, Enterprise (Functioning System).
-
Columns (Interrogatives): What (Data), How (Function), Where (Network), Who (People), When (Time), Why (Motivation).
-
Application to Data: The "What" (Data) column defines the data entities, their relationships, and definitions from each stakeholder's perspective. For example, the "Owner" row in "What" column defines the business-essential data entities (e.g., "Customer," "Order").
III. DATA MANAGEMENT & LINEAGE
Data Lineage Tracking in DataOps
Data Lineage is the lifecycle of data: its origins, movements, transformations, and dependencies.
-
Importance:
-
Debugging: Trace root cause of data errors in pipelines.
-
Compliance/Auditing: Prove data handling for regulations (GDPR, CCPA).
-
Impact Analysis: Understand what reports/models break if a source changes.
-
Data Discovery: Understand how datasets are derived.
-
Pattern-Based Lineage
Automatically infers lineage by analyzing code/configuration patterns.
-
Example (SQL): A parser analyzes a SQL query like:
CREATE VIEW sales_summary AS SELECT region, SUM(amount) FROM orders WHERE date > '2023-01-01' GROUP BY region;- Inferred Lineage:
sales_summary<- (orders.amount,orders.region,orders.date) viaSUMandGROUP BYoperations.
- Inferred Lineage:
-
Tools: OpenLineage, Marquez, custom parsers for Spark/Scala/Python code.
Lineage by Data Tagging
Manually or semi-automatically tags datasets with metadata describing their source and transformations.
-
Example: A dataset
sales_daily_aggis tagged with:-
source_tables: ["orders", "products"] -
transformation: "daily aggregation, joined with product category" -
pipeline_job: "daily_sales_etl_v2"
-
-
Advantage: Works for non-code-based transformations (e.g., ETL tools like Informatica, data catalog entries).
-
Challenge: Requires discipline; can become outdated.
Schema Migration
The process of changing the structure (schema) of a database table or dataset (e.g., adding a column, changing data type).
-
Need: Evolve applications, add new features, improve performance.
-
Example (SQL):
-- Adding a new column to an 'users' table ALTER TABLE users ADD COLUMN phone_number VARCHAR(20); -
Strategies for Zero Downtime:
-
Backward-Compatible Change: Add new columns as nullable. Deploy application that writes to both old and new schema (dual-write). Read from new schema only when ready. Finally, drop old column.
-
Dual-Write: Application writes to both old and new schema structures during transition period.
-
Expand-Contract Pattern: 1) Expand (add new column, backfill), 2) Migrate (switch reads/writes), 3) Contract (remove old column).
-
IV. DATA PROCESSING & INTEGRATION
Extract, Transform, Load (ETL)
A traditional pattern for moving data from source systems to a data warehouse.
Example: Retail Database to Data Warehouse
-
Extract: Connect to operational PostgreSQL
retail_db. Extractorders,order_items,productstables incrementally (using timestamp or CDC). -
Transform:
-
Cleanse: Fix null
customer_id, standardizeproduct_namecase. -
Integrate: Join
orders+order_items+productsto create a denormalizedsales_facttable. -
Aggregate: Calculate daily sales totals per store.
-
-
Load: Write the transformed
sales_factand dimension tables (dim_date,dim_product) into the target Data Warehouse (e.g., Snowflake). UseINSERTfor new records,UPDATE/MERGEfor changes.
[!TIP] Modern Variant: ELT (Extract, Load, Transform) is now common, where raw data is loaded into the warehouse first, then transformed using the warehouse's powerful SQL engine.
Data Warehouse Design: University Setup (Step-by-Step)
Goal: Integrate student records, grades, and finance for reporting.
-
Identify Business Processes: Student Enrollment, Grade Management, Fee Payment.
-
Define Grain: "One line item per student per course per semester" for fact tables.
-
Design Dimensions (Conformed):
-
Dim_Student(student_id, name, program, enrollment_date) -
Dim_Course(course_id, code, title, credits) -
Dim_Time(date_key, day, month, semester, year) -
Dim_Staff(staff_id, name, department) -- for instructors
-
-
Design Fact Tables:
-
Fact_Enrollment(student_key, course_key, time_key, staff_key, grade, credits_earned) -
Fact_Fee_Transaction(student_key, time_key, fee_type, amount, payment_status)
-
-
Choose Schema: Star Schema (fact tables linked directly to dimension tables) for simplicity and query performance.
-
ETL Process: Extract from Student Information System (SIS) and Finance DB, transform to conformed dimensions, load into warehouse.
Real-Time & Time-Series Processing
Processing Real-Time Data with Lambda Architecture
-
Batch Layer: Hourly/daily batch job (Spark) processes all historical clickstream logs from HDFS to compute accurate user session aggregates (e.g., total time on site). Outputs to batch view (e.g., HBase table
batch_user_sessions). -
Speed Layer: Real-time stream (Kafka) of live clicks processed by Flink to compute recent session activity (last 5 mins). Updates real-time view (e.g., Redis cache
realtime_user_sessions). -
Serving Layer: A query service (API) reads from both
batch_user_sessions(HBase) andrealtime_user_sessions(Redis), merges them (e.g., batch data + last 5 min real-time delta), and serves to dashboard.
Detecting Patterns in Time-Series Data (Streaming)
Use a stream processing engine (Flink/Spark Streaming) with windowed operations.
-
Scenario: Detect sudden spike in server CPU usage.
-
Steps:
-
Ingest Stream: Metrics from servers (timestamp, cpu_usage) via Kafka.
-
Window: Apply a tumbling window of 1 minute.
-
Aggregate: Compute average and standard deviation of
cpu_usageper server per window. -
Pattern Detection: Compare current window's average to a baseline (e.g., 30-day moving average from batch layer). If
(current_avg > baseline_avg * 1.5), flag as anomaly. -
Output: Send anomaly alert to monitoring system (e.g., PagerDuty) and store in
anomaliestable.
-
V. TOOLS, PROTOCOLS & AUTOMATION
Webhooks
An event-driven HTTP callback where one application sends an HTTP POST request to another application's URL when a specific event occurs.
-
How it works:
-
Subscriber registers a URL (endpoint) with the Publisher.
-
Event happens in Publisher (e.g., GitHub push, payment success).
-
Publisher serializes event data (JSON payload) and makes an HTTP POST to the Subscriber's URL.
-
Subscriber receives, validates, and processes the payload.
-
-
Example (GitHub): On
pushevent to repomyorg/myrepo, GitHub sends a POST tohttps://my-ci-server.com/github-hookwith payload containing commit details. CI server then triggers a build. -
Payload Example (Payment):
{ "event": "payment.succeeded", "payload": { "invoice_id": "inv_123", "amount": 2500, "currency": "USD" } }
Web Scraping
Automated extraction of data from websites.
-
Techniques:
-
HTML Parsing: Use libraries (BeautifulSoup, Scrapy) to parse DOM.
-
API Reverse-Engineering: Mimic browser's XHR/Fetch calls to get JSON data.
-
Headless Browsers: Control Chrome/Firefox via Puppeteer/Selenium for JS-heavy sites.
-
-
Scenario: rgpvonline.com
-
Goal: Find all representatives with press releases containing "data".
-
Steps:
-
Crawl
https://www.rgpvonline.comto find links to press release pages (e.g.,/press-releases/). -
For each press release page, extract title and content.
-
Filter pages where
titleORcontentcontains the word "data" (case-insensitive). -
Extract representative name (e.g., from byline or metadata).
-
Output: List of
(representative_name, press_release_url).
-
-
-
Ethical/Legal Considerations:
-
Check
robots.txtfor scraping permissions. -
Respect
rate-limiting(don't hammer the server). -
Review Terms of Service; some sites prohibit scraping.
-
Do not scrape personal/sensitive data without consent (GDPR implications).
-
Secure Copy Protocol (SCP)
A network protocol based on SSH for secure file transfer between a local and remote host or between two remote hosts.
-
How it works: Uses SSH for authentication and encryption. Files are copied in binary mode.
-
Command Example:
scp local_file.txt user@remote_host:/remote/directory/ -
Comparison:
-
vs FTP: SCP is encrypted (FTP sends passwords/data in plaintext). SCP is simpler but less feature-rich.
-
vs SFTP: SFTP (SSH File Transfer Protocol) is a subsystem of SSH offering more features (resume, directory listing, remove). SCP is older and simpler; SFTP is more robust and recommended for new systems.
-
Logging, Monitoring, and Alerting (LMA)
A triad for system observability and reliability.
| Component | Purpose | Examples |
|---|---|---|
| Logging | Record discrete events (errors, transactions) with timestamps. | Application logs (JSON format), system logs (/var/log/syslog). |
| Monitoring | Track metrics (CPU, memory, ETL job duration) over time. | Prometheus metrics, CloudWatch dashboards. |
| Alerting | Trigger notifications when metrics/logs breach thresholds. | "Alert if ETL job fails > 3 times in 1 hour" (PagerDuty). "Alert if API latency > 1s" (Slack). |
-
Example for ETL Job:
-
Log: Job writes
"INFO: Started loading table X","ERROR: Connection timeout". -
Monitor: Prometheus scrapes job duration metric
etl_job_duration_seconds. -
Alert: Alertmanager rule:
if etl_job_success == 0 for 10m then page on-call engineer.
-
VI. ADVANCED TOPICS & APPLICATIONS
Data Lake vs. Data Warehouse vs. Data Mart
| Feature | Data Lake | Data Warehouse | Data Mart |
|---|---|---|---|
| Data | Raw, all formats (structured, unstructured). | Processed, structured, modeled. | Subset of warehouse, focused on a domain. |
| Schema | Schema-on-Read. | Schema-on-Write (rigid). | Schema-on-Write. |
| Users | Data scientists, ML engineers. | Business analysts, executives. | Specific business unit (e.g., Sales). |
| Cost | Very low (object storage). | High (compute + licensed software). | Moderate. |
| Agility | High (flexible). | Low (rigid schema changes). | Medium. |
| Use Case | Exploratory analytics, ML, archive. | Standardized reporting, BI. | Department-specific reporting. |
Data Quality & Validation in Pipelines
-
Techniques:
-
Profiling: Check stats (min, max, null count, distinct count) on ingestion.
-
Validation Rules: Define expectations (e.g.,
order_amount > 0,emailmatches regex). Use tools like Great Expectations. -
Reconciliation: Compare record counts and sums between source and target.
-
Anomaly Detection: Monitor data distribution shifts (e.g., sudden drop in daily users).
-
Cloud Data Services (Overview)
-
AWS: S3 (Data Lake), Redshift (Warehouse), Glue (ETL), Kinesis (Streaming).
-
GCP: Cloud Storage (Lake), BigQuery (Serverless Warehouse), Dataflow (Stream/Batch), Pub/Sub (Messaging).
-
Azure: Blob Storage (Lake), Synapse Analytics (Warehouse), Data Factory (ETL), Stream Analytics.
Security in Data Engineering
-
Encryption: At-rest (S3 SSE, TDE) and in-transit (TLS/SSL).
-
Access Control: IAM roles/policies, column-level security, row-level security (RLS).
-
Auditing: Log all data access (CloudTrail, audit logs) and data movement.
-
Linked to Governance: Security policies are enforced via data governance roles (Data Owner defines "who can see PII"). SCP can be used for secure transfer of sensitive data between on-prem and cloud.
[!TIP] Exam Synthesis: Always connect security to governance (e.g., "Data Steward defines access policies; Data Custodian implements via IAM and encryption").