⚡ Category 01 • Core Compute & Ingestion

Engineering: Python, SQL, PySpark & REST APIs

The core programming and compute engine of modern data platforms. From object-oriented Python scripting and high-performance analytical SQL to distributed PySpark transformations and resilient REST API ingestion pipelines.

1. Track Overview & Engineering Foundation

Foundational Pillar

Data Engineering requires bridging procedural scripting, declarative data manipulation, distributed cluster computing, and networked API integration. This track provides end-to-end depth across all four pillars: writing maintainable, type-safe Python packages; constructing sub-second windowed analytical SQL queries; tuning PySpark DAG execution and partition memory; and architecting rate-limited, fault-tolerant REST API extractors.

Python for DE

Iterators, generators, context managers, async I/O, custom decorators, data validation with Pydantic, and modular packaging.

Analytical SQL

Window functions (LEAD, LAG, DENSE_RANK, NTILE), Recursive CTEs, Gaps & Islands, EXPLAIN execution plans, and index tuning.

Distributed PySpark

Catalyst query optimization, Tungsten binary processing, broadcast hash joins, key salting, and partition pruning.

Resilient REST APIs

Sliding window rate limiters, token-bucket throttling, exponential backoff with jitter, idempotent idempotency keys, and pagination.

Topic 1.1

Python Fundamentals, OOP & Advanced Memory Patterns

Intermediate to Advanced

Python is the lingua franca of data engineering. Beyond basic scripting, enterprise data engineering demands generator pipelines to stream gigabyte-scale files without OOM crashes, custom context managers for database transaction isolation, and clean Object-Oriented design patterns.

Generators & Memory Efficiency

Using yield to produce lazy iterables. Reading 100GB logs chunk by chunk with constant O(1) memory overhead.

Context Managers & Scopes

Implementing __enter__ and __exit__ to guarantee database connection closure, S3 socket cleanup, and atomic locks.

Decorators & Cross-Cutting Concerns

Writing parameterised decorators for automated function retry, latency instrumentation, and schema validation.

Type Hinting & Pydantic Contracts

Enforcing static type contracts with typing.TypedDict, Optional, Union, and runtime validation models.

📄 python_chunked_streamer.py
import time
import functools
import logging
from typing import Generator, Dict, Any, Callable

logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("DE_Pipeline")

def retry_with_backoff(retries: int = 3, backoff_in_seconds: float = 1.0):
    """Decorator for resilient I/O operations with exponential backoff and jitter."""
    def decorator(func: Callable):
        @functools.wraps(func)
        def wrapper(*args, **kwargs):
            delay = backoff_in_seconds
            for attempt in range(1, retries + 1):
                try:
                    return func(*args, **kwargs)
                except Exception as e:
                    logger.warning(f"Attempt {attempt}/{retries} failed for {func.__name__}: {e}")
                    if attempt == retries:
                        raise
                    time.sleep(delay)
                    delay *= 2.0
        return wrapper
    return decorator

def stream_large_dataset(file_path: str, chunk_size: int = 10000) -> Generator[list[str], None, None]:
    """Lazy streaming generator keeping memory strictly bounded at O(chunk_size)."""
    batch = []
    with open(file_path, "r", encoding="utf-8") as f:
        # Skip header
        header = f.readline().strip().split(",")
        for line in f:
            batch.append(line.strip())
            if len(batch) >= chunk_size:
                yield batch
                batch = []
        if batch:
            yield batch

# Production usage demonstration
if __name__ == "__main__":
    logger.info("Initialized memory-bounded ingestion generator.")

Production Use Case: 50GB Custodian Settlement Processing

At XXX portfolio analytics, daily custodian transaction feeds arrived as uncompressed 50GB CSV files. Standard pandas.read_csv() failed with out-of-memory errors on worker nodes. By adopting lazy chunked generators coupled with typing validation, memory consumption dropped from 32GB to 140MB with zero file truncation.

Topic 1.2

Advanced SQL, Window Functions & Query Optimization

Advanced

SQL is the language of business logic and analytical serving. A senior data engineer must master analytical window functions, gaps & islands problems, recursive hierarchy traversals, and EXPLAIN ANALYZE execution cost trees.

Analytical Window Functions

Sliding aggregation frames (ROWS BETWEEN 29 PRECEDING AND CURRENT ROW), running cumulative totals, and ranking without group collapse.

Gaps & Islands Technique

Detecting contiguous clusters of activity (e.g. trading streak days, server uptime periods) using difference between row numbers and date intervals.

Common Table Expressions (CTEs) & Recursion

Modularising query logic with CTEs and executing recursive queries to traverse parent-child account or organization hierarchies.

EXPLAIN Plan Analysis & SARGability

Identifying Sequential Scans vs Index Scans, Hash Joins vs Nested Loops, and ensuring predicates remain Search Argument Able (SARGable).

📄 portfolio_peak_drawdown_window.sql
-- Compute Daily Portfolio NAV, High-Water Mark and Maximum Drawdown
WITH daily_nav AS (
    SELECT 
        portfolio_id,
        trade_date,
        SUM(shares_held * closing_price) AS market_value,
        SUM(cash_balance) AS cash_value,
        SUM(shares_held * closing_price + cash_balance) AS total_nav
    FROM investment_positions
    WHERE trade_date >= '2026-01-01'
    GROUP BY portfolio_id, trade_date
),
ranked_metrics AS (
    SELECT 
        portfolio_id,
        trade_date,
        total_nav,
        -- Running Peak-to-Date (High-Water Mark)
        MAX(total_nav) OVER (
            PARTITION BY portfolio_id 
            ORDER BY trade_date 
            ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
        ) AS high_water_mark,
        -- Prior day NAV for daily return calculation
        LAG(total_nav, 1) OVER (
            PARTITION BY portfolio_id 
            ORDER BY trade_date
        ) AS prior_day_nav
    FROM daily_nav
)
SELECT 
    portfolio_id,
    trade_date,
    total_nav,
    high_water_mark,
    ROUND(((total_nav - prior_day_nav) / NULLIF(prior_day_nav, 0)) * 100, 4) AS daily_return_pct,
    -- Drawdown percentage from the all-time peak
    ROUND(((total_nav - high_water_mark) / NULLIF(high_water_mark, 0)) * 100, 2) AS drawdown_pct
FROM ranked_metrics
ORDER BY portfolio_id, trade_date;

Production Use Case: Query Tuning for 200M Row Analytical Dashboard

An executive risk dashboard was timing out after 45 seconds running on a 200M row financial trades table. By replacing non-SARGable WHERE YEAR(trade_date) = 2026 with bounded date ranges WHERE trade_date >= '2026-01-01' AND trade_date < '2027-01-01' and adding a composite index on (portfolio_id, trade_date) INCLUDE (total_nav), execution dropped from 45.2s to 180ms.

Topic 1.3

PySpark & Distributed Cluster Processing

Enterprise Core

Apache Spark provides distributed in-memory computing across hundreds of worker nodes. Mastering PySpark requires deep knowledge of Catalyst optimization rules, Tungsten bytecode generation, wide vs narrow transformations, and data skew mitigation.

Catalyst Optimizer & Execution Plans

Understanding Analysis → Logical Plan → Optimized Logical Plan → Physical Plan. Predicate pushdown, projection pruning, and constant folding.

Wide vs Narrow Transformations

Narrow (map, filter) execute in executor memory without shuffle. Wide (groupBy, join, distinct) force expensive network disk spills and cluster shuffles.

Broadcast Hash Join (BHJ)

Broadcasting dimension tables (< 100MB) to all executors using broadcast(dim_df), eliminating 100% of network shuffle for high-speed joins.

Key Salting for Skew Mitigation

Preventing single-task stragglers where 99% of data lands on one partition by prepending random integer salts [0..N] to the join keys.

📄 pyspark_salting_and_broadcast.py
from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = SparkSession.builder \
    .appName("ProductionEngineeringTuning") \
    .config("spark.sql.adaptive.enabled", "true") \
    .config("spark.sql.adaptive.skewJoin.enabled", "true") \
    .getOrCreate()

def join_skewed_streams(transactions_df, merchants_df, salt_buckets: int = 16):
    """
    Eliminates straggler tasks when merchant_id has massive key skew 
    (e.g., Amazon, Apple accounting for 40% of all volume).
    """
    # 1. Salt the heavy skewed transaction stream
    salted_transactions = transactions_df.withColumn(
        "salt", F.floor(F.rand() * salt_buckets)
    ).withColumn(
        "salted_key", F.concat(F.col("merchant_id"), F.lit("_"), F.col("salt"))
    )

    # 2. Replicate merchant dimension rows across all salt buckets
    replicated_merchants = merchants_df.withColumn(
        "salt", F.explode(F.array([F.lit(i) for i in range(salt_buckets)]))
    ).withColumn(
        "salted_key", F.concat(F.col("merchant_id"), F.lit("_"), F.col("salt"))
    )

    # 3. Join on the salted composite key
    result_df = salted_transactions.join(
        replicated_merchants,
        on="salted_key",
        how="inner"
    ).drop("salt", "salted_key")

    return result_df

Production Use Case: Telecom Churn Pipeline 12-Hour Runtime Reduced to 45 Mins

At XXX Broadband Churn Analytics, the PySpark feature engineering stage hung indefinitely on task 199/200 due to key skew on customer region codes. Applying systematic key salting and enabling Adaptive Query Execution (AQE) distributed the data evenly across 40 executor cores, shrinking pipeline runtime from 12 hours to 45 minutes.

Topic 1.4

Production REST APIs, Pagination & Rate Limit Governance

Enterprise Integration

Most raw data originates from third-party or internal REST/HTTP APIs. An enterprise-grade ingestion system must gracefully handle HTTP 429 Too Many Requests, cursor-based pagination, network drops, and schema drift.

Sliding Window Rate Limiting

Tracking request counts per minute using Redis or local token-buckets to prevent breaching vendor quota limits.

Cursor-Based Pagination

Navigating large datasets using opaque cursor tokens or next_url pointers instead of vulnerable offset/page parameters.

Idempotent Ingestion

Generating deterministic SHA256 checksums from response payloads to guarantee exactly-once persistence even during retries.

Dead Letter Queue (DLQ) Fallback

Quarantining malformed responses or 5xx server errors into S3 DLQ prefixes with SNS alerting without stopping the pipeline.

📄 resilient_api_ingestor.py
import time
import requests
import hashlib
from typing import Dict, Any, Generator

class ResilientAPIExtractor:
    def __init__(self, base_url: str, api_key: str, requests_per_minute: int = 60):
        self.base_url = base_url
        self.headers = {"Authorization": f"Bearer {api_key}", "Accept": "application/json"}
        self.min_interval = 60.0 / requests_per_minute
        self.last_request_time = 0.0

    def _rate_limit(self):
        elapsed = time.time() - self.last_request_time
        if elapsed < self.min_interval:
            time.sleep(self.min_interval - elapsed)
        self.last_request_time = time.time()

    def fetch_paginated_records(self, endpoint: str) -> Generator[Dict[str, Any], None, None]:
        url = f"{self.base_url}/{endpoint}"
        cursor = None

        while url:
            self._rate_limit()
            params = {"cursor": cursor} if cursor else {}
            
            try:
                response = requests.get(url, headers=self.headers, params=params, timeout=10)
                if response.status_code == 429:
                    retry_after = int(response.headers.get("Retry-After", 10))
                    time.sleep(retry_after)
                    continue

                response.raise_for_status()
                data = response.json()

                for record in data.get("items", []):
                    # Compute deterministic idempotency hash
                    record_hash = hashlib.sha256(str(record).encode()).hexdigest()
                    record["_ingest_hash"] = record_hash
                    yield record

                cursor = data.get("next_cursor")
                if not cursor:
                    break

            except requests.RequestException as err:
                # Log and route to DLQ
                print(f"API Extraction Error: {err}")
                break

2. Production Engineering Pipeline Architecture

End-to-End Blueprint

The diagram below illustrates how Python API ingestion, SQL analytical staging, and distributed PySpark clusters interact in a production enterprise deployment.

[REST APIs • WebSockets] → (Python Resilient Ingestor + Rate Limiter) ↓ [AWS S3 Bronze Lake] (Raw JSON / Parquet • Immutable) ↓ [Apache Spark / PySpark Cluster] • Catalyst Optimizer • Key Salting • Broadcast Joins ↓ [AWS S3 Silver Lake] (Cleaned, Reconciled • Delta Lake / Iceberg) ↓ [Analytical SQL Warehouse / PostgreSQL] • Window Functions • NAV & Drawdowns • Power BI

3. Production Best Practices & Performance Tuning

Operational Excellence

1. PySpark Shuffle Partition Sizing

Set spark.sql.shuffle.partitions dynamically based on input size. Target 100MB–200MB per partition instead of the default 200.

2. SARGable SQL Predicates

Never apply functions to indexed columns in WHERE clauses (e.g. WHERE DATE(created_at) = '2026-10-01'). Use timestamp bounds.

3. Python Type Annotations

Always enforce Python 3.12 type annotations with mypy and ruff in CI/CD before any DAG or script deployment.

4. API Circuit Breakers

Implement automated circuit breakers that pause downstream ingestion if upstream error rate exceeds 5% across a 2-minute rolling window.

Hands-On Lab

Lab 01: Build an End-to-End Ingestion & PySpark Lakehouse Pipeline

Lab Guide

In this hands-on lab, you will build a complete end-to-end Python + PySpark pipeline extracting live market data from a REST API, persisting raw JSON to S3 Bronze, and using PySpark to produce a Curated Silver Parquet dataset with rolling 30-day volatility.

1

Set Up Python Ingestion Script

Create ingest_market_api.py with exponential backoff and rate limiting to download hourly OHLCV ticks from Polygon or Alpha Vantage.

2

Write to S3 Bronze Storage

Partition raw JSON by /year=YYYY/month=MM/day=DD/asset_class=EQUITY/ using AWS boto3 with server-side KMS encryption.

3

PySpark Transformations & Quality Checks

Read Bronze JSON, validate schemas with Great Expectations, compute VWAP using window functions, and write to Silver Delta Lake with ACID merge.

4

Verification & Analytical Query

Run SQL queries on the curated Silver layer to verify zero duplicate trades and accurate 30-day rolling moving averages.

4. Senior Data Engineer Interview Preparation

Interview Scenarios
Q1: How do you identify and resolve data skew in PySpark without using broadcast joins? ▲

Answer: First, examine the Spark UI Event Timeline and Summary Metrics. If the median task duration is 2 seconds but the max task duration is 25 minutes with large shuffle read/spill, you have data skew. To resolve it without broadcasting (e.g. both tables are large):
1. Key Salting: Prepend a random integer [0..N-1] to the skewed join key in Table A, and explode Table B with an array of all N integers to match.
2. Adaptive Query Execution (AQE): Set spark.sql.adaptive.skewJoin.enabled = true. Spark automatically splits skewed partitions into smaller sub-partitions at runtime.
3. Filter & Union: Isolate the top 5 skewed keys, process them via a separate high-parallelism branch, and UNION ALL with the unskewed portion.

Q2: Explain how SQL window frames (ROWS vs RANGE) differ in execution and performance. ▼

Answer: ROWS treats each physical record individually based on exact row offsets (e.g., ROWS BETWEEN 1 PRECEDING AND CURRENT ROW). RANGE treats duplicate values in the ORDER BY clause as peers, aggregating all matching values together. In SQL engines, ROWS is significantly faster because it operates on a simple array index without requiring equality checks across rows.

Q3: How do you design an idempotent REST API ingestion pipeline that survives mid-batch network crashes? ▼

Answer:
1. Compute a deterministic content hash (e.g. SHA-256 of the natural business key + event timestamp).
2. Persist the ingested records using an upsert or MERGE INTO operation on the primary/hash key in the target database or Delta Lake table.
3. Maintain a persistent checkpoint table storing the last successfully committed cursor/timestamp.
4. If a job fails at 50%, the retry re-extracts the batch, but the storage engine overwrites matching hashes without creating duplicate rows.

5. Related Specialization Tracks

Explore Next

Ready to practice these concepts live?

Test your SQL and PySpark knowledge with our interactive browser-based runner and 48+ company interview question bank.

▶ Start Practicing Problems 🌊 Next Track: Streaming →
Open Python, SQL & PySpark interview questions →