ML & Data Science

Research survey · LLM systems

Big Data for LLMs: 2024–25 Survey

MSc Data Science · Big Data coursework (Task 5) · Coventry University

LLM progress is no longer just about bigger Transformers. It now depends on the data pipelines, evaluation systems, and serving infrastructure built around them.

Role
Researcher
Year
2025
Focus
Research survey
Stack
9 technologies
  • Spark
  • Ray
  • Daft
  • Data-Juicer
  • vLLM
  • SGLang
  • Flink
  • Vector databases
  • GraphRAG
  • 45

    Cited papers and reports

  • 100T+

    Tokens in RedPajama-V2

  • 0.6%

    Disk needed by LSHBloom vs MinHash

  • 4.48×

    Goodput gain from DistServe

  • 75%

    More requests served with Mooncake

  • 70%+

    Lower job completion time with Hermes

Research from 2024 and 2025 shows a shift from model-centric to system-centric AI. The bottleneck is now engineering: curating trillions of tokens, proving benchmark scores aren't memorised, and serving models fast and cheaply on limited GPU memory. Traditional big data ideas (ETL, stream processing, distributed caching) are being reinvented for these workloads.

The survey follows the LLM lifecycle through three connected phases (Development, Verification, and Deployment) plus how LLMs integrate into real-time pipelines. It draws on recent papers and industry reports from NeurIPS, OSDI, FAST, SIGMOD, VLDB, ACL, and NAACL.

Development covers the data supply chain: open, transparent corpora such as RedPajama-V2 (100+ trillion tokens with 40+ quality signals per document) and Dolma (3 trillion tokens), the move from Spark to Ray, Data-Juicer, and Daft for multimodal processing, and deduplication at scale.

Verification covers data contamination, where benchmark questions leak into training data, and the statistical tests and dynamic benchmarks used to detect it. Deployment covers the inference memory wall and the systems that break through it, and Integration covers streaming, vector search, and knowledge graphs for retrieval-augmented generation.

The LLM lifecycle through a big data lens

  1. 01

    Development

    data supply chain

    Open corpora with provenance (RedPajama, Dolma), multimodal ETL on Ray and Daft, and deduplication with MinHash LSH, SoftDedup and LSHBloom.

  2. 02

    Verification

    contamination

    Separate memorisation from generalisation with ConStat, CDD, and Min-K% Prob, plus dynamic benchmarks like Clean-Eval and LLM-as-a-judge pipelines.

  3. 03

    Deployment

    inference at scale

    KV cache virtualisation (vLLM), prefix caching (SGLang), prefill/decode disaggregation (DistServe, Mooncake), and probabilistic scheduling (Hermes).

  4. 04

    Integration

    real-time RAG

    Flink with ML inference operators, VectraFlow for vector streams, disk-based vector indexes (DiskANN), and GraphRAG with knowledge graphs.

Data processing frameworks

Engines for LLM data pipelines

Apache Spark
JVM executors, excellent for tabular ETL and a mature ecosystem, but high serialisation overhead and GPU contention for model-based filtering.
Ray Data
Python actors with native scheduling across mixed CPUs and GPUs. Strong for heterogeneous, AI-centric workloads.
Daft
Streaming execution engine (Swordfish) that avoids the multimodal 'memory explosion'. Reported 2–7× faster than Ray Data and 4–18× faster than Spark on multimodal tasks.
Data-Juicer
One-stop LLM data-recipe system (SIGMOD 2024, 2.0 in 2025) that dispatches each operator to Ray or Spark and adds a feedback-driven sandbox.

Inference architectures

How each system attacks the memory wall

vLLM · PagedAttention
Splits the KV cache into pages with a page table, like OS virtual memory. Near-zero fragmentation and 2–4× throughput.
SGLang · RadixAttention
Treats the KV cache as a radix tree of reusable prefixes. Up to 6.4× throughput for agentic and few-shot workloads.
DistServe
Runs compute-bound prefill and memory-bound decode on separate GPUs. Up to 4.48× goodput within latency SLOs.
Mooncake
KVCache-centric design (FAST 2025 Best Paper) using cluster CPU RAM and SSDs as a distributed KV store. Serves 75% more requests for Kimi.
Hermes
Models demand as a probabilistic graph and schedules with the Gittins index, cutting average job completion time by over 70%.

Contamination detection

Is a high score reasoning or memorisation?

N-gram overlap
String matching against training data. Precise, but impossible for proprietary data and fooled by paraphrasing.
ConStat
Compares performance on the original benchmark against a rephrased reference. Detects the symptom, not text overlap, and flagged models like Mistral-7B and Llama-3-70B.
CDD
Measures how 'peaked' the output distribution is by sampling. Works on black-box APIs, but sampling is expensive.
Min-K% Prob
Checks the likelihood of the least-probable tokens in a passage. Catches subtle leakage, but needs log-probabilities.
Clean-Eval
Paraphrases and back-translates benchmarks, then verifies semantics with a BERT detector. A proactive defence that shows real drops on clean sets.

Real-time RAG and vector search

Where streaming meets LLMs

Flink + inference
ML inference operators inside dataflows, with async I/O so slow model calls don't stall the stream.
VectraFlow
Stream engine for embeddings (CIDR 2025) with V-TopK similarity search over sliding windows and semantic V-Join.
DiskANN · Vamana
Graph index designed to live on SSD, enabling billion-scale vector search on one node.
GaussDB-Vector · HAKES
Hybrid semantic + relational search inside a DBMS, and a filter-and-refine design for fresh indexes under concurrent writes.
GraphRAG
Retrieves knowledge-graph subgraphs (e.g. from Neo4j) as explicit context to reduce hallucinations.
  • RedPajama-V2 replaces fixed 'hard' filters with 40+ stored quality signals per document, turning the dataset into a queryable database. Ablations showed aggressive filtering can hurt long-tail knowledge.
  • Dolma (AI2, ACL 2024) runs language ID, quality filters, and PII removal across thousands of CPU cores, and tracks licensing and provenance for web-scraped data.
  • Spark struggles with model-based filtering such as BERT toxicity or CLIP alignment scoring: JVM serialisation and GPU contention push teams to Ray and Daft.
  • Deduplication is the most expensive step. MinHash LSH plus distributed Union-Find finds near-duplicate clusters, SoftDedup down-weights common content instead of deleting it, and LSHBloom matches MinHash accuracy with 0.6% of the disk.
  • Contamination threatens every benchmark: a model that saw GSM8K or HumanEval during training is recalling, not reasoning.
  • Black-box methods infer contamination from model behaviour (performance gaps, output peakedness, token likelihoods) because training data is often private.
  • Dynamic evaluation is the strongest defence: generate fresh, semantically equivalent test sets so leaderboard phrasing can't be memorised.
  • LLMOps platforms now stream production interactions through judge models to watch for drift, toxicity, and hallucination spikes.
  • The distinction between big data systems (Spark, Flink, Kafka) and AI systems (training, inference) is dissolving.
  • The edge now comes from infrastructure: transparent data supply chains, forensic evaluation, and inference engines that treat GPU memory like an operating system treats RAM.
  • Next-generation systems such as VectraFlow and Hermes treat data, probability, and vectors as first-class primitives.