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

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

UNIT 2: DATA ENGINEERING - SHORT NOTES

A. FOUNDATIONAL CONCEPTS & ROLES

The Role of a Data Engineer

  • Definition: A data engineer builds and maintains the infrastructure and pipelines that enable data collection, storage, processing, and access for analysis and machine learning.

  • Key Responsibilities:

    • Design, build, and maintain scalable data pipelines (ETL/ELT).

    • Develop and manage data storage solutions (Data Lakes, Warehouses).

    • Ensure data quality, reliability, and security.

    • Collaborate with Data Scientists (provide clean, feature-rich data), Analysts (enable self-service reporting), and Business Stakeholders (translate business needs to technical specs).

  • Real-World Impact: A data engineer at an e-commerce company builds a real-time pipeline to ingest clickstream data, process it for customer behavior analysis, and feed it to a recommendation engine, directly influencing sales and user experience.

[!TIP] Exam Focus: Be prepared to explain the collaborative workflow between Data Engineers, Scientists, and Analysts. Use a concrete example (e.g., fraud detection, supply chain optimization).

The Five V's of Data

V Definition Example
Volume The sheer amount of data generated/collected. Terabytes/Petabytes of IoT sensor readings, social media posts.
Velocity The speed at which data is generated, processed, and absorbed. Real-time stock market ticks, GPS location updates.
Variety The different types and formats of data. Structured (SQL tables), Semi-structured (JSON, XML), Unstructured (images, videos, emails).
Veracity The quality, accuracy, and trustworthiness of data. Incomplete customer records, noisy sensor data, inconsistent formatting.
Value The ultimate usefulness or insight derived from data. Predictive maintenance alerts, personalized marketing offers, trend reports.

[!TIP] Common Pitfall: Don't just list the V's. Always pair each with a specific, relevant example as shown in the table.

Types and Sources of Data

  • By Structure:

    • Structured: Fixed schema, rows & columns (e.g., relational database tables).

    • Semi-structured: No rigid schema but has tags/markers (e.g., JSON, XML, CSV with headers).

    • Unstructured: No predefined model (e.g., text documents, PDFs, images, videos).

  • By Origin:

    • Internal: Generated within the organization (e.g., ERP, CRM, application logs).

    • External: Acquired from outside (e.g., public datasets, social media APIs, purchased data).

  • By Arrival Pattern:

    • Batch: Data collected over a period and processed together (e.g., nightly sales report).

    • Streaming: Continuous flow of data processed in near real-time (e.g., live website analytics).


B. MODERN DATA ARCHITECTURES & STRATEGIES

Modern Data Strategies

  • Cloud-First: Preferring cloud-native services (S3, BigQuery) over on-prem for scalability, cost-efficiency, and managed services.

  • Data Mesh: A decentralized paradigm where data is treated as a product. Domain-oriented teams (e.g., "Finance", "Marketing") own and serve their data as products, with a self-serve data infrastructure platform.

    • Goal: Scale data architectures across large, complex organizations.
  • Data Fabric: An integrated, unified layer of data services and architecture that provides consistent capabilities across multiple hybrid/multi-cloud and on-prem environments. Focuses on abstraction and automation for data access and integration.

    • Goal: Simplify data access and management across disparate sources.

[!TIP] Key Difference: Data Mesh is about organizational structure & ownership (decentralized domains). Data Fabric is about technical integration & abstraction (unified layer).

Modern Data Architecture on Cloud Platforms (Core Components)

A typical pipeline: Ingestion -> Storage -> Processing -> Consumption/Serving.


graph LR

    A[Data Sources<br/>Batch/Stream] --> B[Ingestion Layer<br/>Kafka, Kinesis, Flume];

    B --> C[Storage Layer<br/>Data Lake S3/ADLS<br/>Data Warehouse Redshift/Snowflake];

    C --> D[Processing Layer<br/>Spark, EMR, BigQuery];

    D --> E[Consumption Layer<br/>BI Tools, ML Models, Apps];

  • Cloud Support: Provides elastic scalability (scale up/down compute on demand), managed services (no server/ cluster maintenance), and pay-as-you-go cost models.

Data Lake vs. Data Warehouse

Feature Data Lake Data Warehouse Lakehouse (Modern Convergence)
Purpose Store raw, all data (any format). Store curated, structured data for analytics. Combine both: store raw data + support ACID transactions & BI.
Schema Schema-on-Read (apply schema when reading data). Schema-on-Write (data must conform before loading). Supports both.
Users Data Scientists, Engineers (exploration). Business Analysts, Executives (reporting). Both.
Cost Low-cost object storage (e.g., S3). High-cost specialized compute/storage. Moderate (storage like lake, compute like warehouse).
Examples Amazon S3, Azure Data Lake Storage (ADLS). Amazon Redshift, Snowflake, Google BigQuery. Delta Lake (on S3), Apache Iceberg, Apache Hudi.

[!TIP] Exam Diagram: Be ready to sketch a simple diagram showing Data Lake (raw zone -> curated zone) feeding a Data Warehouse or a Lakehouse architecture.

Creating Scalable & Feasible Infrastructure (Example-Driven)

Process:

  1. Understand Requirements: Volume, velocity, latency needs, user personas (analysts vs scientists), compliance (GDPR, HIPAA).

  2. Choose Storage Paradigm: Start with a Data Lake (S3/ADLS) for flexibility and cost. Add a Cloud Data Warehouse (Snowflake/BigQuery) for fast BI.

  3. Select Processing Engine: Use Spark (via EMR/Databricks) for complex transformations and ML. Use Cloud-Native SQL (BigQuery, Redshift) for simpler ELT.

  4. Orchestrate: Use Airflow or AWS Step Functions to schedule and monitor pipelines.

  5. Secure & Govern: Implement IAM, encryption, and data catalog (AWS Glue, Azure Purview).

  6. Iterate & Scale: Start small (MVP), use auto-scaling groups, monitor costs.

Example (Startup vs. Enterprise):

  • Startup (10 TB, 5 users): S3 (Lake) + Snowflake (Warehouse) + dbt (ELT) + Airflow. Low ops overhead.

  • Enterprise (PB-scale, 1000+ users): Multi-region S3 + Databricks (Unified Spark) + Unity Catalog (Governance) + CI/CD for pipelines. Focus on security, governance, cost allocation.


C. DATA INGESTION & PROCESSING PATTERNS

Data Ingestion: Tools and Patterns

  • Batch Ingestion: Move large volumes of data at scheduled intervals.

    • Tools: Apache Sqoop (RDBMS to HDFS), custom scripts, cloud services (AWS DataSync).
  • Stream/Real-time Ingestion: Continuously capture and process data as it arrives.

    • Tools: Apache Kafka (distributed event streaming platform), Amazon Kinesis (AWS managed), Apache NiFi (dataflow automation).
  • Change Data Capture (CDC): Captures inserts, updates, deletes in a source database in real-time by reading the transaction log.

    • Tools: Debezium (open-source), AWS DMS, Qlik Replicate.

    • Pattern: Source DB -> CDC Tool -> Kafka -> Downstream consumers (Data Lake, Warehouse).

ETL vs. ELT

Aspect ETL (Extract, Transform, Load) ELT (Extract, Load, Transform)
Process Flow Extract -> Transform (in staging) -> Load to target. Extract -> Load (raw) to target -> Transform in target.
Transformation Location Dedicated ETL server/staging area. Inside the target system (Data Warehouse/Lakehouse).
Target System Load Stores only transformed, structured data. Stores raw + transformed data.
Latency Higher (transformation step before load). Lower (load first, transform on-demand).
Tools Traditional: Informatica, Talend. Cloud: AWS Glue (can do both). Modern: dbt (data build tool), BigQuery, Snowflake, Databricks SQL.
Best For Complex, legacy transformations; small, structured sources. Cloud-native, scalable systems; data science exploration; when raw data needs preservation.

[!TIP] Rule of Thumb: If your target is a modern cloud data warehouse/lakehouse (BigQuery, Snowflake, Databricks), ELT is the dominant pattern due to scalability and flexibility.

Big Data Processing Frameworks

Apache Hadoop Ecosystem (Historical Significance)
  • HDFS (Hadoop Distributed File System): Distributed, fault-tolerant storage. Splits files into blocks (default 128MB) stored across cluster nodes.

  • MapReduce: Programming model for processing vast datasets.

    • Map Phase: Processes input key-value pairs to produce intermediate key-value pairs.

    • Shuffle & Sort: System groups intermediate values by key.

    • Reduce Phase: Aggregates values for each key to produce final output.

    • Time Complexity: $O(n)$ for map, $O(k \log k)$ for shuffle/sort (where k is number of unique keys), $O(m)$ for reduce. Overall dominated by I/O.

    • Limitation: High latency due to disk I/O between map and reduce stages.

  • YARN (Yet Another Resource Negotiator): Cluster resource manager. Schedules and manages resources for applications (MapReduce, Spark) on the Hadoop cluster.

  • Diagram:

    DiagramSEARCH: HDFS MapReduce YARN architecture diagram

Apache Spark (Modern Standard)
  • Key Features:

    • In-Memory Processing: Caches data in RAM across iterations, 100x faster than MapReduce for iterative algorithms (ML).

    • DAG Execution: Creates a Directed Acyclic Graph of operations, optimizes the entire workflow (pipeline) before execution.

    • Unified Engine: Single platform for Batch Processing, Stream Processing (Spark Streaming/Structured Streaming), SQL (Spark SQL), ML (MLlib), Graph Processing (GraphX).

    • APIs: Available in Scala, Java, Python (PySpark), R.

  • Advantages over MapReduce: Speed (in-memory), ease of use (high-level APIs), unified stack, better for iterative/streaming workloads.

Amazon EMR (Elastic MapReduce)
  • Definition: A managed cluster platform that simplifies running big data frameworks (Spark, Hadoop, Hive, HBase, Presto) on AWS.

  • Role in Simplifying:

    • Provisioning: Automatically launches and configures EC2 instances.

    • Scaling: Add/remove instances automatically based on workload.

    • Integration: Native integration with S3 (storage), Glue (catalog), Redshift (warehouse).

    • Cost: Pay only for the time the cluster runs. Supports Spot Instances for 60-90% cost savings.

    • Managed: Handles cluster maintenance, patching, and failure recovery.

Specialized Architectures for Real-Time Processing

Lambda Architecture
  • Goal: Balance low-latency (real-time) with accuracy (recomputed from all data).

  • Three Layers:

    1. Batch Layer: Stores immutable, raw master dataset. Computes batch views (accurate, but slow) from all historical data.

    2. Speed Layer: Processes real-time stream to compute real-time views (low-latency, but approximate). Compensates for batch layer latency.

    3. Serving Layer: Merges batch views and real-time views to serve user queries. Provides a complete, up-to-date view.

  • Diagram:

    DiagramCANVAS: Draw three horizontal layers. Batch Layer (HDFS/Data Lake) -> Batch Views (Hive/Spark). Speed Layer (Kafka/Storm) -> Real-time Views (Redis/Druid). Serving Layer (Query Engine) merges both to answer queries.

  • Challenge: Complexity of maintaining two separate codebases (batch & speed).

Kappa Architecture
  • Simplification: Eliminates the batch layer. All data is processed as a single, immutable stream.

  • Process: Ingest all data (historical + real-time) into a single, scalable log (Kafka). Use a stream processor (Spark Streaming, Flink) to compute views continuously. Serving layer stores results from the stream.

  • Benefit: Simpler (one codebase). Requires the stream processor to be powerful enough to reprocess historical data if needed.

  • Diagram:

    DiagramCANVAS: Single stream: Data Sources -> Kafka (central log) -> Stream Processor (Spark/Flink) -> Serving Layer (Cassandra/Elasticsearch).


D. DATA PREPARATION, WRANGLING & GOVERNANCE

Data Wrangling and Data Discovery

  • Process: The iterative process of cleaning, transforming, enriching, and validating raw data into a usable format for analysis.

    1. Discovery: Understanding data structure, quality, and relationships (using profiling tools).

    2. Cleaning: Handling missing values, outliers, duplicates, inconsistent formats.

    3. Transformation: Structuring (pivoting), normalizing, aggregating, joining.

    4. Enrichment: Adding derived columns, joining with external datasets.

    5. Validation: Ensuring data meets quality rules and schemas.

  • Tools at Scale: Apache Spark (DataFrames), dbt (SQL-based transformations), Pandas (for smaller data), cloud services (AWS Glue DataBrew).

Data Governance & Lineage

  • Importance: Ensures data quality, security, compliance, and trust. Defines ownership (who is responsible for a dataset?), policies, and standards.

  • Data Lineage Tracking in DataOps:

    • Technical Lineage: "How" data moves and transforms. Tracks the flow from source -> transformations -> destination at the column/field level.

      • Example: SalesDB.Orders.Order_Date -> (Spark job clean_orders.py) -> DataWarehouse.FactSales.Order_Date.
    • Business Lineage: "What" the data represents from a business perspective. Maps technical assets to business terms, reports, and KPIs.

      • Example: FactSales.Order_Date -> "Daily Revenue Report" -> "Business KPI: Monthly Recurring Revenue".
    • Implementation Methods:

      • Pattern-Based Lineage: Parse job scripts/configurations (SQL, Spark code) to extract data flow. Tool Example: OpenLineage.

      • Lineage by Data Tagging: Attach tags/metadata (source system, PII, owner) to datasets/tables. Lineage inferred from tag propagation. Tool Example: Apache Atlas, Collibra.

Gartner Data Maturity Model

A framework to assess an organization's data management maturity. Stages:

  1. Unaware: No formal data management. Data is siloed, inconsistent.

  2. Emerging: Awareness grows. Basic tools (spreadsheets, reports) used. Limited governance.

  3. Consolidating: Centralized data team forms. First data warehouse/lake. Basic data quality rules.

  4. Defined: Formal data governance program. Standardized processes, metadata management, data quality monitoring. Data is a strategic asset.

  5. Optimized: Data is embedded in business processes. Advanced analytics, AI/ML pervasive. Continuous improvement of data products. Data drives competitive advantage.

[!TIP] Example: A company at "Consolidating" might have a single Snowflake instance. At "Defined", they have a data catalog, data stewards, and SLA-defined data quality. At "Optimized", all product decisions are A/B tested using real-time data.

Schema Management

  • Schema Evolution: The process of changing a schema (add/remove/rename columns) over time as data needs change.

  • Schema Migration Strategies:

    • Backward Compatibility: New schema can read old data. Rule: Only add new columns with default values or make fields optional.

    • Forward Compatibility: Old schema can read new data. Rule: Only remove columns that are nullable or have defaults; never change data types in incompatible ways.

  • Schema Registry: A centralized repository for schemas (often for serialized data like Avro, Protobuf). Ensures producers and consumers agree on data format.

    • Example: Confluent Schema Registry with Kafka. Producers register a schema ID; consumers fetch schema by ID to deserialize. Enforces compatibility rules (BACKWARD, FORWARD, FULL).

E. MACHINE LEARNING INTEGRATION & PIPELINES

Machine Learning Lifecycle (Key Stages)

  1. Problem Definition: Frame business problem as ML task (classification, regression).

  2. Data Collection & Preparation: Critical role for Data Engineer. Ingest, clean, and store relevant data.

  3. Preprocessing & Feature Engineering: Most impactful stage for model performance. Create informative features from raw data.

  4. Model Training: Select algorithm, train on prepared dataset. Tune hyperparameters.

  5. Model Evaluation: Assess performance on hold-out test set using metrics (Accuracy, F1, RMSE).

  6. Model Deployment: Package model and deploy as real-time endpoint or batch inference job.

  7. Model Monitoring & Maintenance: Track prediction drift, data drift, performance degradation. Retrain as needed.

Role of the Data Engineer in ML

  • Build and maintain scalable, reproducible data pipelines for:

    • Training: Providing large, clean, versioned training datasets.

    • Inference: Serving features in real-time for online predictions or generating batch predictions.

  • Implement feature stores (centralized repositories of features) to ensure consistency between training and serving.

  • Automate data validation and schema enforcement for ML datasets.

  • Set up ML tracking and metadata logging (e.g., MLflow, SageMaker Experiments).

AWS SageMaker (Cloud ML Service)

  • Components:

    • Notebook Instances: Managed Jupyter notebooks for exploration.

    • Training Jobs: Managed, scalable training on custom data/algorithm. Auto-scales.

    • Hyperparameter Tuning: Automatically runs multiple training jobs with different HPs to find optimal model.

    • Endpoints: Real-time, scalable REST APIs for model inference. Auto-scales based on traffic.

    • Pipelines: CI/CD service for building, versioning, and automating ML workflows (data prep -> training -> deployment).

  • How it Supports Scalability: All components are serverless/managed. You specify resources (instance type, count), AWS handles provisioning, scaling, and load balancing. Integrates with S3, Glue, and other AWS services.

Pre-processing and Feature Engineering

  • Critical Role: "Garbage in, garbage out." Quality features are often more important than complex algorithms. Can increase model accuracy by 10-20%.

  • Common Techniques:

    • Handling missing values (imputation).

    • Encoding categorical variables (One-Hot, Label Encoding).

    • Scaling/normalization (StandardScaler, MinMaxScaler).

    • Creating interaction features, polynomial features.

    • Date/time feature extraction (hour, day of week, holiday flag).

    • Text processing (TF-IDF, embeddings).

  • Engineering Scalable Pipelines: Implement transformations as distributed Spark jobs or within dbt models. Store engineered features in a feature store (e.g., Feast, SageMaker Feature Store) for reuse.


F. OPERATIONAL EXCELLENCE, SECURITY & COMPLIANCE

Securing Cloud Storage and Data

Layered Security Model (Example: AWS S3):

  1. Physical & Network Security: AWS data center security, VPC (private network).

  2. Encryption:

    • In Transit: TLS/SSL for data moving to/from S3.

    • At Rest: Server-Side Encryption (SSE-S3, SSE-KMS, SSE-C) or Client-Side Encryption.

  3. Access Control (IAM & Policies):

    • IAM Policies: Attached to users/roles, define allowed actions (s3:GetObject).

    • Bucket Policies: Attached to S3 bucket, resource-based.

    • Object ACLs: Legacy, per-object permissions (use sparingly).

  4. Network Controls: VPC Endpoints (S3 gateway endpoint) allows private access from VPC without traversing the public internet.

  5. Monitoring & Logging: S3 Access Logs, AWS CloudTrail for API calls.

Diagram:

DiagramSEARCH: AWS S3 security model diagram showing layers: User -> IAM Policy -> VPC Endpoint -> Bucket Policy -> Object Encryption.

Logging, Monitoring, and Alerting

  • Importance: Ensure reliability, performance, and cost control of data pipelines. Detect failures early.

  • Key Metrics to Track:

    • Pipeline Health: Job success/failure rate, duration, throughput (records/sec).

    • Data Quality: Row counts, null percentages, freshness (time since last update), schema drift.

    • Resource Utilization: CPU/Memory of clusters, disk I/O, network I/O.

    • Cost: Cost per job, cost per TB processed.

  • Tools:

    • Cloud-Native: Amazon CloudWatch Logs/Metrics, Azure Monitor, GCP Cloud Monitoring.

    • Open Source: Prometheus (metrics collection) + Grafana (dashboards), ELK Stack (Elasticsearch, Logstash, Kibana for log analysis).

CI/CD for Data Pipelines

Applying DevOps practices to data infrastructure:

  1. Infrastructure as Code (IaC): Define cloud resources (clusters, storage) in code (Terraform, AWS CloudFormation). Enables versioning, reproducibility.

  2. Pipeline as Code: Define data workflows in code (Airflow DAGs, dbt projects, SageMaker Pipelines).

  3. Automated Testing:

    • Unit Tests: Test individual transformation functions.

    • Integration Tests: Test pipeline end-to-end on a small, isolated dataset.

    • Data Quality Tests: Assert expectations on output data (using Great Expectations, dbt tests).

  4. Deployment Automation: Use CI/CD tools (GitHub Actions, Jenkins, GitLab CI) to automatically run tests and deploy pipeline code/IaC to dev/staging/prod environments.

  5. Monitoring in Prod: Automated alerts on pipeline failures, data quality breaches, cost anomalies.

Data Governance Frameworks

  • The Zachman Framework: A schema for organizing enterprise architecture artifacts (models, documents) from multiple perspectives (Planner, Owner, Designer, Builder, Subcontractor, Functioning Enterprise) and aspects (What, How, Where, Who, When, Why). Used to ensure all architectural concerns are addressed comprehensively.

  • Policies for Privacy & Compliance:

    • Data Privacy: Implement data masking/anonymization (PII removal), obtain consent, manage data subject requests (right to be forgotten).

    • Compliance: Design systems to meet regulations:

      • GDPR (EU): Strict rules on personal data, consent, breach notification.

      • HIPAA (US Healthcare): Safeguards for Protected Health Information (PHI).

    • Auditability: Maintain immutable logs of all data access and transformations for forensic audits.


G. SPECIALIZED TOOLS & TECHNIQUES (SHORT NOTES)

Apache Hadoop (Architecture and Core Components)

  • Core Idea: Distributed storage and processing of very large datasets across clusters of commodity hardware.

  • HDFS Architecture:

    • NameNode: Master server. Manages file system namespace and metadata (file->block mapping). Single Point of Failure (mitigated by HA).

    • DataNode: Slave node. Stores actual data blocks. Reports block status to NameNode.

    • Block: Default 128MB. Replicated (default 3x) across DataNodes for fault tolerance.

  • MapReduce Programming Model:

    • User defines map() and reduce() functions.

    • Framework handles splitting, distribution, sorting, fault tolerance.

    • Limitation: High disk I/O between map and reduce stages makes it slow for iterative/real-time tasks.

  • YARN: Separated resource management from the MapReduce programming model. Allows multiple processing engines (Spark, Tez) to share cluster resources.

  • Historical Significance: First widely adopted open-source big data framework. Foundation for the modern data ecosystem. Largely superseded by Spark for processing, but HDFS remains a common storage layer (often replaced by cloud object storage).

Amazon EMR (Managed Hadoop/Spark Service)

  • What: Fully managed cluster service for running Spark, Hadoop, Hive, HBase, Presto, Flink.

  • Key Features:

    • Provisioning: Launch clusters in minutes. Choose instance types (compute/memory optimized).

    • Scaling: Auto-scaling policies. Add/remove instances based on CloudWatch metrics.

    • Cost-Effective: Use Spot Instances for fault-tolerant workloads (up to 90% discount). EBS-backed or instance-store.

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

    • Steps API: Define a series of processing steps (JAR, Spark, Hive, etc.) to run on the cluster. Simplifies workflow.

  • Role: Eliminates undifferentiated heavy lifting of cluster setup, configuration, tuning, and management. Allows focus on application logic.

Modern Data Strategies

  • Data Mesh: Decentralized domain-oriented data ownership. Domains (e.g., "Customer", "Supply Chain") produce and own their data as products. A central self-serve data platform provides tools (storage, pipelines, governance) to all domains.

    • Pros: Scales organizationally, improves domain accountability, faster domain innovation.

    • Cons: Requires significant cultural shift, potential for duplication, needs strong platform team.

  • Data Fabric: Integrated, unified layer of data services (catalog, governance, integration, quality) that provides consistent capabilities across hybrid/multi-cloud and on-prem environments. Uses metadata, AI/ML, and knowledge graphs to automate data discovery, integration, and preparation.

    • Pros: Simplifies access, improves agility, handles complexity of multiple sources.

    • Cons: Can be complex to implement initially, vendor-dependent.

  • Lakehouse Architecture: Merges the flexibility and cost of a Data Lake (raw data on S3) with the management, ACID transactions, and performance of a Data Warehouse (via table format like Delta Lake, Apache Iceberg, Apache Hudi on top of the lake).

    • Pros: Single source of truth, supports both BI & ML, avoids data duplication, cost-effective.

    • Cons: Relatively new, ecosystem still maturing, requires careful management of table formats.

CI/CD for Data/ML

  • Goal: Automate the testing, building, and deployment of data pipelines and ML models to ensure reliability, speed, and reproducibility.

  • Key Practices:

    1. Version Control: All code (pipelines, models, IaC) in Git.

    2. Automated Testing:

      • Data/Code Tests: Unit tests for transformations, data quality assertions (Great Expectations, dbt tests).

      • Model Validation: Performance regression tests on hold-out datasets.

    3. Automated Deployment: CI/CD pipeline (GitHub Actions, Jenkins) triggers on merge to main:

      • Run all tests.

      • Build pipeline artifacts (dbt, Airflow DAGs).

      • Deploy to staging environment.

      • Run integration/smoke tests.

      • Approve and promote to production (canary/blue-green).

    4. Infrastructure as Code (IaC): Terraform/CloudFormation for reproducible environments.

    5. Monitoring & Rollback: Monitor deployed pipelines/models. Automated rollback on failure.

  • Tools: dbt (built-in testing/deployment), MLflow (model packaging/tracking), Kubeflow/SageMaker Pipelines (ML), Airflow/Prefect (orchestration), Terraform (IaC).

Streaming Analytics Pipeline (End-to-End Flow)

DiagramCANVAS: Draw a left-to-right flow with numbered stages.

  1. Ingestion: Data Sources (IoT, logs, clicks) -> Message Broker (Apache Kafka/Kinesis). Handles high-throughput, durable storage of events.

  2. Processing: Kafka -> Stream Processor (Apache Spark Structured Streaming/Apache Flink/Kafka Streams). Performs windowing, aggregations, joins, anomaly detection. Outputs to:

    • Serving Layer (for low-latency queries): Redis, Apache Druid, ClickHouse.

    • Storage Layer (for batch/ML): Data Lake (S3/ADLS) in file formats (Parquet/Delta).

  3. Serving/Storage: Processed data lands in:

    • Real-time DB/OLAP: For dashboards (Grafana, Superset).

    • Data Warehouse/Lakehouse: For historical analysis, ML training.

  4. Consumption: BI Tools (Tableau, Power BI), ML Models (real-time inference), Applications (alerts, APIs).

  5. Orchestration & Monitoring: Airflow/Step Functions to manage batch jobs, Prometheus/Grafana for pipeline metrics, CloudWatch for logs.

Secure Copy Protocol (SCP)

  • Definition: A command-line utility based on SSH for secure transfer of files between a local and a remote host or between two remote hosts.

  • How it Works: Uses SSH for authentication and encryption. Provides confidentiality and integrity of data in transit.

  • Syntax: scp [options] source_file user@remote_host:destination_path

  • Use Case in Data Engineering: Securely moving sensitive data files (e.g., PII, credentials) between on-prem servers and cloud VMs, or between EC2 instances. Superseded for large-scale/batch transfers by tools like rsync (over SSH) or cloud-native services (AWS DataSync, Azure Data Factory), but remains a fundamental secure transfer tool.

Webhooks

  • Definition: A user-defined HTTP callback (or small code snippet) that is triggered by a specific event in a source system. It's a way for one application to push real-time data to another application.

  • How it Works:

    1. User configures a webhook URL (endpoint) in Application A (e.g., GitHub, Stripe).

    2. When event occurs (e.g., "code pushed", "payment succeeded"), Application A sends an HTTP POST request with event data (JSON) to the configured URL.

    3. Application B (the receiver) processes the payload (e.g., triggers a CI/CD pipeline, updates a database).

  • Role in Event-Driven Architectures: Enables loose coupling and real-time reactivity. Alternative to polling. Core mechanism for integrating SaaS services and triggering data pipeline steps.

  • Example: A GitHub webhook posts to an Airflow webhook endpoint when a PR is merged, triggering a data pipeline rebuild.

Web Scraping

  • Definition: The automated extraction of data from websites using software.

  • Techniques & Tools:

    • HTML Parsing: Use libraries like BeautifulSoup (Python) or Cheerio (Node.js) to parse HTML/XML and extract data using CSS selectors/XPath.

    • Dynamic Content: Use headless browsers like Selenium or Puppeteer to render JavaScript-heavy sites.

    • Frameworks: Scrapy (Python) - full framework for large-scale crawling and scraping.

  • Ethical Considerations & Legality:

    • Check robots.txt: Respects site's scraping policy.

    • Review Terms of Service: Often prohibits scraping.

    • Don't Overwhelm Servers: Implement rate limiting, use delays.

    • Copyright & Data Ownership: Scraped data may be copyrighted. Personal data scraping may violate privacy laws (GDPR, CCPA).

    • Public vs. Private Data: Scraping public websites is generally legal but ethically nuanced. Scraping behind logins is often a violation.

  • Example (Hypothetical): To find representatives with press releases about "data":

    1. Crawl https://www.rgpvonline.com (hypothetical government site).

    2. For each representative's page, scrape links to press releases.

    3. For each press release, fetch content and search for keyword "data".

    4. Output list of representative names where keyword found.

    • Code Snippet (BeautifulSoup Concept):

      
      soup = BeautifulSoup(html_content, 'html.parser')
      
      press_links = soup.select('.press-release a')  # CSS selector
      
      for link in press_links:
      
          if "data" in link.text.lower():
      
              print(representative_name)
      
      
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