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 PillarData 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.
Python Fundamentals, OOP & Advanced Memory Patterns
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.
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.
Advanced SQL, Window Functions & Query Optimization
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).
-- 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.
PySpark & Distributed Cluster Processing
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.
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.
Production REST APIs, Pagination & Rate Limit Governance
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.
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 BlueprintThe diagram below illustrates how Python API ingestion, SQL analytical staging, and distributed PySpark clusters interact in a production enterprise deployment.
3. Production Best Practices & Performance Tuning
Operational Excellence1. 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.
Lab 01: Build an End-to-End Ingestion & PySpark Lakehouse Pipeline
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.
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.
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.
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.
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
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.
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.
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🌊 Streaming
Kafka, Spark Streaming, event-driven pipelines, and exactly-once processing.
Go to Streaming →📈 Financial Data
Market data, ETF weightings, portfolio analytics, and custodian reconciliation.
Go to Financial Data →🏛️ Data Platforms
Delta Lake, Snowflake, Kimball dimensional modeling, and PostgreSQL.
Go to Data Platforms →☁️ Cloud
AWS, Azure, S3 lakehouse architecture, and Azure Databricks.
Go to Cloud →