Skip to content
CY-703 (C) · Data Engineering/Quick Revision Short Notes

Data Engineering (CY-703 (C)) - Unit 5 Short Notes

UNIT 5: Data Engineering


I. Foundations and Lifecycle

Data Engineering: Definition & Scope

  • Definition: The practice of designing, building, and maintaining the systems and infrastructure that enable the collection, storage, processing, and serving of data for analytical and operational use.

  • Evolution: From traditional ETL in data warehouses to modern cloud-native, scalable, and real-time data platforms (Lakehouse, Data Mesh).

  • Lifecycle Core Components:

    1. Ingestion: Acquiring data from source systems.

    2. Storage: Persisting data in appropriate systems (lakes, warehouses).

    3. Processing: Transforming, cleansing, and preparing data.

    4. Serving: Making data available to consumers via APIs, dashboards, etc.

    5. Consumption: End-users (analysts, data scientists, apps) using the data.

Role of Data Engineer

  • Responsibilities: Build and maintain data pipelines, ensure data quality and reliability, optimize performance/cost, implement security/governance.

  • Skills: Programming (Python, SQL, Scala), cloud platforms (AWS/Azure/GCP), big data frameworks (Spark, Hadoop), orchestration (Airflow), data modeling.

  • Impact: Enables data-driven decisions by providing trusted, accessible data.

Five V's of Data

V Definition Example
Volume Scale of data (TB, PB, EB) Social media feeds, IoT sensor data
Velocity Speed of data generation/processing Real-time stock trades, clickstreams
Variety Different formats & structures Text, images, JSON, CSV, video
Veracity Data quality, accuracy, trustworthiness Incomplete records, inconsistent formats
Value Business insights derived from data Predictive maintenance, customer segmentation

Types & Sources of Data

  • Structured: Fixed schema, relational DBs (e.g., MySQL tables).

  • Unstructured: No predefined schema (e.g., images, PDFs, videos).

  • Semi-structured: Flexible schema (e.g., JSON, XML, logs).

  • Sources: Internal (CRM, ERP) vs. External (APIs, public datasets, web scraping).

  • Batch vs. Stream: Batch (large volumes, scheduled, e.g., daily sales reports) vs. Stream (continuous, low-latency, e.g., fraud detection).

Data-Driven Decision Making

  • Process: Define objective → Collect relevant data → Process & analyze → Generate insights → Make decision → Monitor outcome.

  • Example: Retailer analyzes real-time sales + inventory + weather data to dynamically adjust pricing and stock levels (demand forecasting).

[!TIP] Exam Focus: Be ready to give specific examples for each "V" and clearly differentiate batch/stream ingestion with use cases.


II. Modern Data Architecture

Key Architectural Paradigms

Architecture Core Idea Pros Cons
Lakehouse Combines data lake flexibility with warehouse structure/management Cost-effective, supports ML & BI, ACID transactions Relatively new, tooling evolving
Data Mesh Decentralized domain-oriented ownership (data as a product) Scalable, domain expertise, reduces bottlenecks Complex governance, cultural shift
Data Fabric Unified layer providing integrated access across disparate sources Simplifies access, abstracts complexity Implementation overhead, vendor lock-in risk

Cloud Platforms (Key Services)

  • AWS: S3 (Storage), Redshift (Warehouse), Glue (ETL), Kinesis (Streaming), EMR (Spark/Hadoop).

  • Azure: Blob Storage, Synapse Analytics, Data Factory, Stream Analytics, Databricks.

  • GCP: Cloud Storage, BigQuery, Dataflow (Apache Beam), Pub/Sub, Dataproc.

Diagram Concept:

DiagramCANVAS: Cloud Data Platform - Show S3/Blob/Storage as central raw zone, with Glue/Dataflow processing into Redshift/BigQuery/Synapse, serving to BI tools and ML services.

Data Lake vs. Data Warehouse

Feature Data Lake Data Warehouse
Schema Schema-on-read (flexible) Schema-on-write (rigid)
Data Raw, all formats (structured/unstructured) Processed, structured only
Users Data scientists, ML engineers Business analysts, executives
Cost Low-cost object storage Higher-cost, optimized compute
Use Case Exploratory analytics, data science Reporting, BI, dashboards
  • Hybrid Approach: Use lake for raw storage, warehouse for curated business data (Lakehouse pattern).

Scalable Infrastructure Design (Example: E-commerce)

  1. Assess: Volume (10M users), Velocity (real-time clicks), Variety (logs, transactions, images).

  2. Choose Cloud: AWS for breadth of services.

  3. Design:

    • Ingestion: Kinesis (stream), S3 (batch uploads).

    • Storage: S3 (data lake - raw/processed zones), Redshift (curated warehouse).

    • Processing: Glue (ETL), EMR (Spark for large-scale ML).

    • Serving: Redshift Spectrum (query S3), QuickSight (BI).

  4. Optimize: Auto-scaling, spot instances for EMR, S3 lifecycle policies.

Modern Data Strategies

  • Business Alignment: Start with business outcomes (e.g., "reduce churn by 5%").

  • Agility: Modular pipelines, CI/CD for data, cloud-native services.

  • Cost Optimization: Pay-as-you-go, storage tiering, right-sizing compute, serverless options (Glue, BigQuery).


III. Data Ingestion and Integration

Batch vs. Stream Ingestion

Aspect Batch Ingestion Stream Ingestion
Data Size Large volumes Small, continuous records
Latency High (minutes to days) Low (milliseconds to seconds)
Processing Scheduled (e.g., nightly) Continuous
Tools Sqoop, custom scripts, batch Glue jobs Kafka, Kinesis, Pub/Sub
Use Case Daily sales summary Real-time fraud detection

Purpose-Built Ingestion Tools

  • Change Data Capture (CDC): Captures row-level changes in DBs (e.g., Debezium, AWS DMS). Used for near-real-time replication.

  • APIs: REST/gRPC endpoints for structured data pull/push.

  • Messaging Systems: Kafka, RabbitMQ for decoupled, durable event streaming.

Webhooks

  • Concept: User-defined HTTP callbacks triggered by an event (reverse API).

  • Workflow: Event occurs → Source system sends HTTP POST to pre-configured URL (webhook endpoint) → Your service processes payload.

  • Example: GitHub webhook sends event data to your CI/CD server on every git push.

Web Scraping

  • Techniques: HTML parsing (BeautifulSoup, Scrapy), headless browsers (Selenium) for JS-heavy sites.

  • Scenario (Regulatory Data): Scrape government websites (rgpvonline.com) for press releases containing "data". Use Python requests + BeautifulSoup, filter by keyword, store in S3.

    
    # Pseudo-code
    
    response = requests.get(url)
    
    soup = BeautifulSoup(response.text, 'html.parser')
    
    articles = soup.find_all('div', class_='press-release')
    
    for article in articles:
    
        if "data" in article.text.lower():
    
            save_to_s3(article)
    
    
  • Considerations: robots.txt, rate limiting, legal compliance (terms of service).

CI/CD in Data Pipelines

  • Automation: Git for code, automated testing (data quality checks, schema validation), orchestrated deployment (Airflow, Prefect).

  • Tools: Jenkins, GitLab CI, GitHub Actions, dbt Cloud.

  • Process: Code commit → Run tests (Great Expectations) → Deploy to staging → Validate → Promote to prod.

[!TIP] Common Pitfall: Confusing webhooks (server push) with APIs (client pull). Web scraping must respect robots.txt and terms of use to avoid legal issues.


IV. Storage Solutions

Data Lake Storage Patterns & Zones

  • Zones:

    1. Raw (Bronze): Immutable copy of source data.

    2. Processed (Silver): Cleansed, conformed, validated data.

    3. Curated (Gold): Business-level aggregates, ready for consumption.

  • Merits: Low cost, schema flexibility, supports all data types, central repository.

  • Applications: Data science sandbox, regulatory archive, IoT data hub.

  • Example: Raw JSON logs → Silver Parquet with schema → Gold aggregated daily reports.

Data Warehouse Design: Dimensional Modeling

  • Goal: Optimize for query performance and business understanding.

  • Star Schema: One Fact Table (metrics, foreign keys) connected to multiple Dimension Tables (descriptive attributes).

    • Example (Retail): fact_sales (product_id, store_id, date_id, revenue) linked to dim_product, dim_store, dim_date.
  • Snowflake Schema: Normalized dimensions (dimensions have sub-dimensions). Reduces redundancy but adds join complexity.

  • University Setup Example:

    • Fact: fact_enrollment (student_id, course_id, semester_id, grade).

    • Dimensions: dim_student (name, program, year), dim_course (title, credits, dept), dim_semester (term, year).

Cloud Storage Security

  • Encryption:

    • At Rest: Server-side (SSE-S3/KMS), client-side.

    • In Transit: TLS/SSL.

  • IAM & Access Controls: Principle of least privilege. Use IAM roles/policies, bucket policies, ACLs (minimize).

  • Secure Copy Protocol (SCP): Command-line tool using SSH for secure file transfer between local & remote host (alternative to insecure FTP).

    
    scp -i mykey.pem file.txt ec2-user@ec2-ip:/path/
    
    
  • Diagram Concept:

    DiagramCANVAS: Cloud Storage Security Layers - Show data encrypted at rest in S3 bucket, IAM policy attached to user/role, bucket policy allowing only specific VPC endpoint, data in transit via HTTPS/SSL.


V. Data Processing Paradigms

ETL vs. ELT

Aspect ETL (Extract-Transform-Load) ELT (Extract-Load-Transform)
Sequence Transform before loading into target Load raw data first, transform in target
Target Optimized warehouse (schema-on-write) Data lake/cloud warehouse (schema-on-read)
Transformation Done in staging/ETL server Pushed to powerful target (e.g., Snowflake, BigQuery)
When to Use Legacy warehouses, complex cleansing needed Cloud data warehouses/lakes, fast ingestion, iterative exploration
Tools Informatica, Talend, SSIS dbt, Spark SQL, BigQuery ML

ETL Process (Retail Example)

  1. Extract: Pull orders, order_items, products from transactional MySQL DB.

  2. Transform:

    • Clean: Handle nulls, correct data types.

    • Join: Merge tables on keys.

    • Aggregate: Calculate total_sales per product/day.

    • Enrich: Add product category from dim_product.

  3. Load: Insert into fact_daily_sales table in Redshift.

Lambda Architecture (Batch + Speed Layers)

  • Goal: Balance latency & accuracy for big data.

  • Layers:

    1. Batch Layer: Processes all historical data (Hadoop/Spark). Creates batch view (accurate, high latency).

    2. Speed Layer: Processes real-time stream (Spark Streaming, Flink). Creates real-time view (low latency, approximate).

    3. Serving Layer: Merges batch & real-time views to serve queries (e.g., via API).

  • Diagram Concept:

    DiagramCANVAS: Lambda Architecture - Show incoming data flowing to both Batch Layer (HDFS -> Spark -> Batch View) and Speed Layer (Kafka -> Streaming -> Real-time View). Serving Layer queries combine both views for output.

Kappa Architecture (Simplified Real-Time)

  • Idea: Treat all data as stream. Use a single stream processing engine (e.g., Apache Flink) for both real-time and historical reprocessing.

  • Simplification: Eliminates batch layer. To recompute history, replay stream from持久化 log (Kafka).

  • When: When batch processing logic is a subset of stream logic.

Streaming Analytics Pipeline

  • Components: Source (Kafka) → Stream Processor (Flink/Spark Streaming) → Sink (Data Warehouse/NoSQL/Dashboard).

  • Tools: Kafka (ingestion), Flink/Spark Streaming (processing), Cassandra/Redis (sink), Grafana (visualization).

  • Use Cases: Real-time dashboards, anomaly detection, alerting.

Scaling Considerations for Stream Processing

  • Throughput: Records/sec. Scale by increasing parallelism (partitions, executors).

  • Latency: Time from event to result. Minimize by windowing, state management optimization.

  • Fault Tolerance: Checkpointing (to HDFS/S3), exactly-once semantics, replication (Kafka).

  • State Management: Handle large state ( RocksDB ), savepoints for recovery.

[!TIP] Key Difference: Lambda uses two codebases (batch & speed), Kappa uses one. ELT leverages modern cloud warehouse compute; ETL does work upstream.


VI. Big Data Processing Frameworks

Apache Hadoop

  • HDFS: Distributed file system. Blocks (128MB), replication (3x), NameNode (metadata), DataNodes (storage).

  • MapReduce: Programming model. map() (filter/sort) → shuffle → reduce() (aggregate). Disk-based, high latency.

  • YARN: Resource manager. Separates resource management from scheduling.

  • Architecture Diagram:

    DiagramSEARCH: Hadoop HDFS architecture diagram showing NameNode, DataNodes, blocks, and MapReduce flow.

Apache Spark

  • Key Features:

    • In-Memory Processing: Caches data in RAM → 100x faster than MapReduce for iterative algorithms.

    • DAG Scheduler: Creates logical execution plan, optimizes workflow.

    • APIs: RDDs (low-level), DataFrames/Datasets (high-level, optimized).

    • Unified Engine: Batch, Streaming (Spark Streaming), ML (MLlib), Graph (GraphX).

  • Advantages over MapReduce:

    • Speed (in-memory vs. disk).

    • Ease of use (high-level APIs in Python/Scala/Java/R).

    • Unified stack for diverse workloads.

    • Better iterative processing (ML, graph).

Amazon EMR

  • Managed Service: Runs Hadoop/Spark clusters on AWS.

  • Simplifies: Cluster provisioning, scaling, configuration, monitoring.

  • Integration: Native with S3 (storage), IAM (security), CloudWatch (monitoring).

  • Use Case: Large-scale ETL, log processing, ML on petabyte-scale data without managing infrastructure.


VII. Machine Learning Pipeline Integration

Key Stages of ML

  1. Data Collection & Storage: Ingest raw data into lake/warehouse.

  2. Data Preprocessing & Feature Engineering: Clean, transform, create features. Most time-consuming, highest impact on performance.

  3. Modeling: Select algorithm, train on training set.

  4. Evaluation: Test on holdout set, metrics (accuracy, F1, RMSE).

  5. Deployment: Serve model as API (real-time) or batch prediction.

  6. Monitoring & Retraining: Track drift, performance decay, retrain.

Pre-processing & Feature Engineering

  • Importance: "Garbage in, garbage out." Quality features > complex algorithms.

  • Techniques:

    • Handling missing values (impute, drop).

    • Encoding categoricals (one-hot, label).

    • Scaling/normalization (StandardScaler, MinMax).

    • Creating interaction features, date parts, text vectorization (TF-IDF).

  • Impact: Can improve model accuracy by 10-30%+.

AWS SageMaker

  • End-to-End ML Development: Built-in Jupyter notebooks, built-in algorithms, automatic model tuning, one-click deployment.

  • Components:

    • Studio: IDE for ML.

    • Training: Managed training jobs on any algorithm (bring your own or built-in).

    • Processing: Data processing/feature engineering jobs.

    • Deployment: Real-time endpoints (REST API) or batch transform jobs.

    • Monitoring: Detect drift, track metrics.

  • Scalability: Automatically scales compute for training/deployment; integrates with S3 for data.

ML Infrastructure on AWS (Diagram Concept)

DiagramCANVAS: AWS ML Infrastructure - Show S3 as central data lake. SageMaker Processing job reads from S3, writes features to Feature Store. SageMaker Training job reads from Feature Store/S3, outputs model to S3. SageMaker Endpoint deploys model from S3. Model serves predictions to application. CloudWatch monitors everything.


VIII. Data Governance, Security, and Operations

Data Governance

  • Policies: Rules for data quality, security, privacy (GDPR, HIPAA).

  • Ownership & Stewardship: Data owners (business) define policies; data stewards (technical) implement.

  • Implementation Framework: People + Processes + Technology (catalogs, lineage tools, policy engines).

Data Lineage Tracking

  • Pattern-Based: Automatically infers lineage by analyzing pipeline code (e.g., dbt manifest.json, Airflow DAGs).

  • Lineage by Data Tagging: Manually or automatically tag datasets with source/transformation metadata; build graph from tags.

  • DataOps Context: Critical for impact analysis, debugging, compliance, understanding data flow.

  • Example: If sales_report table is wrong, lineage shows it depends on raw_orders → staging_orders → fact_orders. Pinpoint issue source.

Data Maturity Models

  • Gartner: 5 levels (Awareness → Reactive → Defined → Managed → Optimized). Assess across people, process, technology.

  • Zachman Framework: 6x6 matrix (What, How, Where, Who, When, Why) x (Planner, Owner, Designer, Builder, Subcontractor, Functioning Enterprise). Focuses on enterprise architecture, not just data.

  • Assessment & Roadmap: Evaluate current state, identify gaps, prioritize initiatives (e.g., "implement data catalog first").

Logging, Monitoring, and Alerting

  • Tools: Cloud-native (CloudWatch, Azure Monitor), open-source (Prometheus, Grafana), ELK stack.

  • Metrics: Pipeline success/failure rates, data freshness, record counts, latency, resource utilization (CPU, memory).

  • Incident Response: Alerts (Slack, PagerDuty) → Runbook → Triage → Fix → Post-mortem.

  • Example: Airflow DAG failure → CloudWatch alarm → PagerDuty page → Engineer checks logs, fixes dependency, reruns.

Secure Cloud Storage Practices

  • Encryption: Always enable at-rest (SSE-KMS) and in-transit (HTTPS).

  • Key Management: Use cloud KMS (AWS KMS, Azure Key Vault) for customer-managed keys (CMK). Rotate keys.

  • Access Control: IAM roles over users, bucket policies restricting by VPC/prefix, no public access unless required.

  • Compliance: Enable logging (S3 access logs, CloudTrail), use AWS Config rules for compliance as code.

  • Diagram Concept:

    DiagramCANVAS: Secure S3 Bucket - Show bucket with SSE-KMS encryption, bucket policy allowing access only from specific IAM role within a VPC endpoint, CloudTrail logging all API calls to a separate audit bucket.


IX. Advanced Topics and Specialized Concepts

Schema Migration

  • Process: Analyze source/target schemas → Design mapping → Develop transformation logic → Test with subset → Execute (often dual-write) → Validate → Cutover.

  • Challenges: Data type mismatches, null handling, referential integrity, downtime minimization.

  • Tools: Custom scripts, AWS DMS (supports heterogeneous migration), Apache Sqoop.

  • Example: Migrating from on-premise Oracle (VARCHAR2) to cloud PostgreSQL (TEXT). Map data types, handle CLOB/BLOB conversion.

Pattern Detection in Time-Series Streaming Data

  • Techniques: Sliding window aggregations, anomaly detection (statistical: Z-score, ML: Isolation Forest), trend/seasonality detection.

  • Architecture: Stream processor (Flink) maintains stateful windows → Apply detection logic → Output alerts to Kafka topic → Dashboard/Alerting system.

  • Example: Detect spike in server CPU usage (>3σ above 1-hour moving average) within 5-minute window.

Data Wrangling & Data Discovery

  • Data Wrangling (Cleaning/Transformation): Handling missing values, outlier treatment, data type conversion, normalization. Tools: OpenRefine, dbt, Pandas.

  • Data Discovery (Cataloging/Profiling): Automated scanning of data assets to extract schema, statistics (min/max, distinct count), sample values, PII detection. Tools: AWS Glue Data Catalog, Azure Data Catalog, Amundsen.

  • Goal: Make data understandable and trustworthy for analysts.

Real-Time Processing with Lambda Architecture (Step-by-Step)

  1. Batch Layer: Hourly Spark job on S3 raw data → computes accurate aggregates → writes to batch view (Redshift).

  2. Speed Layer: Flink job consumes Kafka clickstreams → computes recent 5-min counts → writes to real-time view (Redis).

  3. Serving Layer: API queries Redis for real-time count + Redshift for historical → merges results → serves to dashboard.

  4. Reconciliation: Periodically (daily) batch job corrects any inaccuracies in real-time view.

[!TIP] Exam Focus: Be able to draw and explain Lambda/Kappa architectures. For governance, contrast pattern-based vs. tagging lineage. Schema migration requires a clear, step-by-step plan.

Go to where you left off?

Quick Add to Notes

Save questions, your own notes and screenshots into notes filed by unit. It takes a free account.

Create free account

Have an account? Log in