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.
- Researcher
- 2025
- Research survey
- 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
The problem
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.
Overview
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.
Structure
The LLM lifecycle through a big data lens
- 01
Development
Open corpora with provenance (RedPajama, Dolma), multimodal ETL on Ray and Daft, and deduplication with MinHash LSH, SoftDedup and LSHBloom.
- 02
Verification
Separate memorisation from generalisation with ConStat, CDD, and Min-K% Prob, plus dynamic benchmarks like Clean-Eval and LLM-as-a-judge pipelines.
- 03
Deployment
KV cache virtualisation (vLLM), prefix caching (SGLang), prefill/decode disaggregation (DistServe, Mooncake), and probabilistic scheduling (Hermes).
- 04
Integration
Flink with ML inference operators, VectraFlow for vector streams, disk-based vector indexes (DiskANN), and GraphRAG with knowledge graphs.
Tools & tech
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.
Development: the data supply chain
- 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.
Verification: trustworthy evaluation
- 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.
Conclusion
- 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.