📋 KEY INSIGHTS
- Big data is defined not just by volume but by the “3 Vs”: Volume (terabytes to petabytes), Velocity (real-time or near-real-time data streams), and Variety (structured, semi-structured, and unstructured data from diverse sources).
- Apache Spark is the dominant distributed processing framework for batch and streaming data — it processes data in memory across a cluster, making it 10–100x faster than Hadoop MapReduce for iterative algorithms like machine learning.
- Apache Kafka is the standard for real-time event streaming — it acts as a durable, fault-tolerant message bus that decouples data producers from consumers and can handle millions of events per second.
- The modern data lakehouse architecture combines the storage flexibility of a data lake (raw files in S3 or GCS) with the ACID transactions and query performance of a data warehouse, using open table formats like Apache Iceberg or Delta Lake.
- Cloud-native data warehouses (Snowflake, BigQuery, Redshift) have largely replaced on-premises Hadoop clusters for analytical workloads — they separate compute from storage, enabling pay-per-query pricing and elastic scaling.
- For data scientists, understanding the big data ecosystem matters primarily for feature engineering at scale, training ML models on datasets that don’t fit in memory, and real-time model serving with streaming features.
The majority of data science work happens on datasets that comfortably fit in a single machine’s memory. But as organisations collect more data — clickstreams, sensor readings, transaction logs, user-generated content — the point where standard pandas and scikit-learn workflows break down is reached sooner than most teams expect. A 10GB CSV already strains a laptop. A year of website clickstream data at a medium-sized company is measured in terabytes. At this scale, the tools change: distributed computing frameworks, streaming pipelines, and cloud data warehouses replace single-machine processing. This guide gives you the conceptual and practical foundations of the big data ecosystem that data scientists work with in 2026: Apache Spark, Apache Kafka, the data lakehouse architecture, and cloud data warehouses.
Apache Spark — Distributed Batch and Stream Processing
Apache Spark is the most widely used framework for large-scale data processing. It extends the MapReduce programming model in two fundamental ways: it processes data in memory rather than writing intermediate results to disk after every step, and it supports a rich set of transformations beyond map and reduce. Spark’s core abstraction is the Resilient Distributed Dataset (RDD) — a fault-tolerant, immutable, partitioned collection of records that can be operated on in parallel across a cluster. In practice, most Spark code today uses the higher-level DataFrame API, which provides SQL-like operations with automatic query optimisation through Spark’s Catalyst query planner and Tungsten execution engine.
Spark’s execution model divides a job into a DAG (Directed Acyclic Graph) of stages. Within each stage, tasks run in parallel on different partitions of the data. Transformations (filter, groupBy, join, select) are lazy — they build up an execution plan but do not execute until an action (collect, count, write, show) is called. This laziness allows Spark to optimise the entire execution plan before running a single task, eliminating redundant operations and choosing efficient join strategies.
| Component | What It Does | Data Scientist Use Case |
|---|---|---|
| Spark Core | Task scheduling, memory management, fault tolerance | Underlying engine for all Spark jobs |
| Spark SQL / DataFrame | Structured query processing, schema inference | Feature engineering on large tabular datasets |
| MLlib | Distributed ML algorithms (linear models, trees, clustering, ALS) | Training models on datasets too large for scikit-learn |
| Spark Structured Streaming | Real-time stream processing with exactly-once semantics | Real-time feature computation for online serving |
| GraphX | Distributed graph computation | Graph analytics, community detection, PageRank |
| Delta Lake | ACID transactions on Spark data | Reliable feature store with time travel and schema enforcement |
The most common performance problems in Spark are data skew (one partition is much larger than others, causing that executor to become the bottleneck while others are idle) and shuffle operations (join and groupBy require shuffling data across the network, which is expensive). Diagnosing performance issues requires reading Spark’s execution plan and the Spark UI — the key metrics to watch are shuffle read/write bytes, executor peak memory, and task duration distribution (high variance in task duration signals skew).
For data scientists, the most practical Spark skills are: reading and writing data in Parquet format (columnar storage, compressed, efficient for analytical queries); using Spark SQL for complex feature engineering that involves window functions, joins across large tables, and aggregations; and using MLlib’s ALS for recommendation systems and its feature engineering tools (VectorAssembler, StringIndexer) as preprocessing for distributed model training. PySpark — Spark’s Python API — allows you to write Spark jobs in Python with pandas-like syntax, with the computation distributed across the cluster.
Apache Kafka — Real-Time Event Streaming
Apache Kafka is a distributed event streaming platform designed for high-throughput, low-latency, durable message passing. In the modern data stack, Kafka plays the role of the central nervous system: data producers (applications, databases, IoT sensors, microservices) publish events to Kafka topics, and data consumers (stream processors, data warehouses, ML inference services, analytics dashboards) subscribe to those topics independently. This decoupling means producers and consumers are developed, scaled, and deployed independently — a fundamental advantage over point-to-point integrations.
Kafka’s key architectural concepts: Topics are named streams of events, analogous to database tables but append-only. Partitions split a topic across multiple brokers for parallel consumption and horizontal scaling. Consumer groups allow multiple consumers to read from a topic in parallel — each partition is read by exactly one consumer in the group, enabling horizontal scaling of consumption. Offsets track the position of each consumer group within each partition, enabling exactly-once semantics and replay from any historical point. Retention keeps messages for a configurable period (default 7 days), enabling consumers to reprocess historical data — Kafka is not just a queue but a persistent, replayable event log.
| Pattern | What It Solves | Example |
|---|---|---|
| Event sourcing | Durable audit log of all state changes | Every database change published as Kafka event |
| Real-time feature computation | Compute ML features from live event streams | Rolling 1-hour purchase count for fraud detection |
| Change Data Capture (CDC) | Sync databases to data lake without bulk exports | Debezium reads PostgreSQL WAL → Kafka → S3 |
| Stream-to-stream join | Enrich events with context from another stream | Join clickstream with user profile stream |
| Fan-out | One event triggers multiple downstream consumers | New order → inventory, analytics, email, fraud check |
| Replay / backfill | Reprocess historical events with new logic | Retrain model on last 30 days of events |
The Data Lakehouse — Architecture for 2026
The data lake was supposed to solve the limitations of the data warehouse: cheap object storage (S3, GCS) could hold any data format at any schema, making it easy to land raw data quickly and define the schema later. But in practice, data lakes became “data swamps” — files with no documentation, incompatible schemas, no ACID guarantees, and query performance far worse than a columnar data warehouse. The data lakehouse architecture, popularised by Databricks and adopted by Snowflake, BigQuery, and DuckDB, attempts to combine the best of both.
The key innovation is open table formats — Apache Iceberg, Delta Lake, and Apache Hudi — that add a metadata layer on top of raw Parquet files. This metadata layer provides ACID transactions (concurrent writes without corruption), time travel (query data as of any past timestamp), schema evolution (add columns without rewriting data), and efficient data skipping (partition pruning and column statistics that allow the query engine to skip irrelevant files). The result is a storage layer that can be queried at near-data-warehouse speed while retaining the raw, unprocessed data in its original format.
| Platform | Type | Storage | Compute Scaling | Best For |
|---|---|---|---|---|
| Snowflake | Cloud data warehouse | Proprietary (on S3/GCS/Azure) | Automatic, per-second billing | BI, analytics, data sharing |
| Google BigQuery | Serverless data warehouse | Proprietary Colossus | Fully serverless | Ad-hoc queries on petabyte scale |
| Amazon Redshift | Cloud data warehouse | Redshift-managed + S3 | Cluster or serverless | AWS-native analytics workloads |
| Databricks Lakehouse | Lakehouse (Delta Lake) | Customer-owned S3/GCS/ADLS | Spark clusters, auto-scaling | ML + analytics on same platform |
| DuckDB | In-process analytical DB | Local files or S3 | Single machine (very fast) | Local analytics on medium data (up to ~500GB) |
| Apache Flink | Stateful stream processor | Kafka / S3 | Cluster-based | Complex real-time streaming pipelines |
For data scientists, the practical implications of the lakehouse architecture are significant. Feature engineering can now happen directly on the raw data lake using SQL or PySpark, without a separate ETL step into a warehouse. Time travel enables reproducible ML training: you can point your training pipeline at the feature table as it existed on a specific date, guaranteeing that retraining the model six months later with the same code produces the same dataset. The feature store — a centralised repository of computed, versioned features — is built on top of a lakehouse table format in most modern ML platforms.
✦ SUMMARIZE THIS ARTICLE WITH AI
ETL pipelines and orchestration with Airflow that feed data into the lakehouse are covered in our ETL Pipelines guide. The data engineering fundamentals — batch vs streaming, data modelling, warehouse concepts — are in our Data Engineering Fundamentals guide. For the SQL skills required to work with Snowflake, BigQuery, and Spark SQL, see our Advanced SQL guide. Data engineering interview questions on Spark, Kafka, and distributed systems are in our Data Engineering Interview Q&A.



