Amazon Web Services (AWS) is a cloud computing platform and subsidiary of Amazon that provides on-demand computing resources and services including compute, storage, databases, networking, analytics, machine learning, and developer tools. AWS offers these services to individuals, enterprises, and governments on a pay-as-you-go basis.
DatologyAI is a data curation service for AI model builders. It processes proprietary enterprise data, public web data, and licensed data for pre-training and mid-training on open models or privately deployed foundation models. Its four-stage pipeline cleans poorly formatted, empty, short, or evaluation-contaminated data; curates it by quality, task, and taxonomy; creates synthetic data for task relevance and diversity; and composes sources into staged training datasets. The company also describes seeded rephrasing methods for generating synthetic data at trillion-token scale.
EKS on HyperPod is an AWS infrastructure service for running Kubernetes-orchestrated H100 GPU clusters. DatologyAI used it as the execution environment for a large-scale synthetic-data pipeline combining Ray, KubeRay, and vLLM, with workflows spanning data curation, synthesis, training, and evaluation.
KubeRay is a platform for running Ray workloads within Kubernetes clusters. In the described DatologyAI pipeline, it was used with Ray and vLLM to orchestrate large-scale synthetic-data jobs on EKS on HyperPod.
Kubernetes (K8s) is an open-source container orchestration system hosted by the Cloud Native Computing Foundation. It manages containerized applications across multiple hosts, providing mechanisms for deploying, maintaining, scaling, and scheduling workloads. Applications can be configured with YAML and deployed across hybrid-cloud environments; the platform is also used for production operations involving storage, databases, networking, and security hardening.
Ray is an open-source, Python-native framework for managing, executing, and optimizing distributed computing workloads, developed by Anyscale. Its core provides tasks, actors, and objects for scaling Python applications, while higher-level libraries support data processing, distributed model training, batch and LLM inference, model serving, reinforcement learning, and end-to-end generative-AI workflows. Ray can coordinate heterogeneous CPU and GPU resources and scale workloads from a laptop to clusters with thousands of GPUs.
SGLang is an open-source serving framework and inference engine for large language models and multimodal models, developed by the SGLang project. It provides an alternative serving stack for deployed models, with runtime and kernel optimizations for inference workloads. Its documented capabilities include RadixAttention, a zero-overhead batch scheduler, cache-aware load balancing, structured-output decoding, speculative decoding, disaggregated prefill and decode, model parallelism, and support for GPU and TPU backends. The project also publishes integrations and optimizations for current open models and multimodal, image, and video-generation workloads.
Slurm is a workload orchestration system used to schedule synthetic-data generation, model-training, and evaluation jobs on a separate GPU cluster.
Spark is a data-curation tool used to run jobs that prepare prompts and documents for synthetic-data generation. In DatologyAI's workflow, it supports the curation and synthesis stages used to produce large-scale training data.
vLLM is an open-source library and inference-serving engine for deploying large language models on GPUs and other hardware accelerators. Originally developed in UC Berkeley's Sky Computing Lab, it is maintained by a broad community and provides an OpenAI-compatible API server, as well as Anthropic Messages API and gRPC support. Its serving stack manages attention key-value memory with PagedAttention, continuous batching, chunked prefill, prefix caching, streaming generation, and disaggregated prefill, decode, and encode. It also supports speculative decoding, quantization, optimized attention and GEMM/MoE kernels, automatic kernel generation and graph-level transformations, structured outputs, tool calling, multiple decoding algorithms, and distributed inference through tensor, pipeline, data, expert, and context parallelism. The project integrates with Hugging Face models and supports decoder-only, mixture-of-experts, hybrid state-space, multimodal, embedding, retrieval, reward, and classification models. It can run across NVIDIA, AMD, Intel, and other supported accelerators, as well as x86, ARM, and PowerPC CPUs, and can be installed with uv or pip or built from source.
Searchable transcript of Lessons from Generating 12 Trillion Synthetic Tokens — Bogdan Gaza, DatologyAI — AI Engineer (20:32). Search for a phrase, then click its timestamp to jump straight to that moment in the video.
Captions sourced from the original video on YouTube, published by AI Engineer. The video, its captions and all related intellectual property remain the property of their respective owners; AINotes claims no ownership. Provided for research, accessibility and search — see the Transcript Notice and Copyright Policy.
00:01 [music] We're doing this. Hey everyone, welcome to my talk. I'm Bogdan. I'm one of the co-founders and the CEO of Theology. Today we're going to talk about the engineering lessons from scaling synthetic data for trillion token scale pre-training. Um, let's see. So why do we need synthetic data in the first place? Well, turns out that with the scaling laws, we kind of hit a data wall um you need to put exponentially more compute and exponentially more data to get um models that are getting kind of linearly better and
00:40 better. So um we're we're at this point where pre-training runs are massive. So in order to do the the kind of the ones that actually get to I don't know the mythos level um capabilities you need a lot of data and a lot of it uh cannot be found on the web. So um internet has limited amount of data um we can supplement it with high quality synthetic data.
01:03 So what are we going to talk about today? We're going to talk about um you know how do we do this in the first place? We're going to talk about a bit about beyond web which is kind of our syntheta data recipe. Um we're going to talk about the engineering lessons from our prior experience here at Ethology of running trillion tail trillion um token scale synthetic data runs.
01:24 Um we recently completed one for about 12 trillion tokens. That's across web math and code. And we're going to talk about what we learned from running absolutely massive synthetic runs. Um and then we're going to talk about kind of what's the why does this matter in the first place. Um awesome. So let's get started. Um, so the way that we solve synthetic data at theology and maybe let's start with a with a quick intro about it.
01:49 We're a kind of data curation service. Think of us as the oil refinery for your kind of pre-training and mid-training data sets. We work with customers both startups and enterprises that come to us to um kind of help them with their pre-training and mid-training data needs. So in general they have a lot of data. um they want to get to the best data sets possible and we help them identify the the highest quality data points and um once we do that we pass them through our synthetic data recipe the ones that we're going
02:16 to talk about today um to get absolutely massive data sets. So um that's a bit about us uh in terms of beyond web is is our synthetic data recipe. Um there's kind of many ways to do synthetic data but um kind of the the way that we approach it is um quite different. So here's a few results from from um kind of how we do this. So uh there's two plots.
02:38 Maybe let's start with the first one. Um on the x-axis we have the model size and billions of parameters. Um 1 billion, 3 billion, 8 billion. And then on the y-axis we have um accuracy. So um we've taken kind of a number of data sets like Neotron Synth, Cosmopedia, QA rap, um Red Pajama, and then we've trained models at kind of 1B, 3B and 8B scale together with um kind of the the recipe we've put together for for uh synthetic data which we call beyond web.
03:07 Um and what we see is that um kind of our synthetic data recipe um kind of matches and and sometimes even um kind of in this case the the 3B data point matches the 8B data point um for neatron synth. So that means that you can train slightly smaller models for the same performance as you would train larger ones by using synthetic data. So um that's a bit about kind of the the the quality of of our recipe.
03:35 Um same thing here when you when you talk about performance. So on the x axis we have um kind of number of tokens in billions and then the y- axis we have uh the the average accuracy. Um and what we show is that we match the the the performance of um kind of the the um Neotron synth and the Cosmopedia ones in about 2.7x and 5.3x respectively uh less tokens.
03:57 So uh if you want to train a model to the same performance, you can do that um by using less data and and overall less compute. Um so at a high level what what beyond web is is is kind of how we think about our synthetic data recipe. What we do is we do rephrasings. Um so where we we kind of see the rephrase the the synthetic data process by um kind of finding the highest quality data points in a customer's data set and then we rephrase or restructure this data points in in various ways.
04:26 Um this is results are actually a bit old. I think the the origin we published this in uh in um I think late summer 25 and at this point we're we're much much better than this um and and we're continuing to improve quite a bit. So how do you get synthetic data right in the in the simplest form you have some sort of target data set again we're doing rephrasings um or reframings of existing documents.
04:49 So the idea behind it is that um if you just prompt a model and say hey give me a synthetic data um data point you're going to just learn the modes of the distribution the data was was origin trained on. So uh the way that we do it as mentioned previously is that we seed it with um some sort of existing documents. So we identify high quality documents in the customer's um training data sets and then we we pass them through through our pipeline.
05:14 At the simplest set, we have a a target data set. Imagine this can be text, this can be image text. Um, in the simplest forms, it is just a bunch of prompts that that we prefix to um a number of documents that we've identified. You have a bunch of GPUs and then you um kind of use various models to generate synthetic data. Um so in practice, you have some sort of um expectations around how much time this is going to take.
05:40 But in practice, this is a lot more wavy and and it it really it's really hard to match the the kind of ideal throughput that you're trying to get of your um out of your GPUs. So, how do we manage that? Well, well, where did we start? We started where um synthetic data wasn't integrated in our product was something that we would do in a research setting uh on a cluster that that runs slurm that is away from our core kubernetes-based pipeline.
06:03 Uh we had this two-track system where kind of the product and the research team would would not use the same type of codebase. This was slow, errorprone, um, limited in velocity and kind of led to underutilizing our GPUs. We have a set number of GPUs like everybody out there. And in general, we want to be able to to maximize the the the use of our GPUs by running this in a in a cohesive way both with training and synthetic data runs.
06:28 Um, so where do we start? In general, um, we'd like to run everything on on Kubernetes. Um, a lot of our pipeline is either Ray or Spark. Today we're going to talk about um kind of ray and cube and vlm um those are kind of the core components of our synthetic data library um or data um kind of pipelines. So today what we do is we have an orchestration layer in one kubernetes cluster.
06:50 We um run spark in another Kubernetes clusters and um we we keep the aster tracker which is our data catalog in the original control cluster where everything um kind of gets coordinated um into um all the data is stored in S3 that's our storage layer um as long as it speaks kind of S3 compatible APIs we can run it. So in some ways we everything needs to run in Kubernetes because our pipelines get deployed in all sorts of environments.
07:16 Um it can be AWS, it can be GCP, it can be onrem etc. We want to make sure that that um kind of we can run it in any um Kubernetes environment. And then separately we had this H100 cluster. Um it ran synthetic, it run training, it ran eval kind of a number of things. All of them originally orchestrated by slurm. Um moving data in between these two clusters was manual.
07:35 So whenever you generated synthetic data in slurm, you had to backfill it into the the asset tracker into our data catalog. So where are we now? So um kind of same setup in some ways where we have the the original cluster um that that does the control. We have the compute cluster that that we run spark jobs and now we have this other H100 compute cluster in in our case we run everything on AWS.
07:55 We run it on hyper pod. Hyperpod has this product called um EKS on hyper pod. So it's basically a bunch of H100 clusters on the same spine all of them um orchestrated by uh by Kubernetes. So what we do is we continue to run spark and um ray jobs on demand capacity on demand CPUs and sometimes on demand GPUs um especially for for things that we don't require H100s we can do that um and then we orchestrate it in a similar way that that we've done before where kind of spark runs within its Kubernetes cluster and then ray
08:24 runs within um the same cluster for um GPUs that are not H100s and anything that we need to do for um H100 nodes we go to the this other Kubernetes cluster and we we kind of have a job scheduler that we've built internally that that orchestrates all of these things. Um if we do this then this allows us to to run fast experiments. Imagine that you can run kind of um curate where you run a bunch of spark jobs, you run a bunch of synthetic data, then you go and and launch a training job and then you run an eval job and
08:54 you can orchestrate all of this together in one workflow which um for our research team is is quite useful and this is puts the the the research flow and the engineering flow in in the same lines in the same infrastructure in the same code in the same path. Okay. So, um we where are we now, right? Like where Oh, look. I have a a oneonone with someone on this fine.
09:15 Um where where we're at is um kind of synthetic data is is a core thing to to our workflows. Um we have we run experiments all the time. We we do this crosscluster scheduling. It's a lot more faster. It's a lot more reliable and we can run it at scale. So, um we've recently done a massive um synthetic data run. Uh let's talk a bit about the the engineering lessons that that we learned from this right.
09:37 So um what are the bottlenecks at at trillion token scale? Um well first and foremost um internet scale um for for text data set is is quite limited right on on the web we can find about 30 trillion tokens give or take for um images there's about 2 12.8 8 billion in in things like dataccom for or for text we usually start from existing highquality data sets like fine web dcmotron um we can we can um kind of put them all together and we get this this large data sets um but at this scale um usually everything gets stored
10:12 in parquet or some sort of of format you have millions of data partitions um if you're trying to fetch all of them from s3 at a high rate you might get throttled and then if you especially if you run this on ray um your your ray metadata uh is becomes and how you manage that becomes a bottleneck and that leads to the the worst thing possible which is that you get idle GPUs.
10:32 So bottleneck um that that we initially seen is that we have um S3 we try to fetch things from our array head. um this has some sort of metadata store. um this is kind of quite slow and and if you do one request that takes 10 milliseconds and and um you're trying to to do I don't know millions and millions of them add I don't know let's say um if you need to do um kind of the metadata fetches originally those are more expensive um and there's also a rate limit so if you need to do millions of them you're going to get
11:05 rate limited then your workers are going to be quite slow um so we kind of ran the the the back of the anvil of math just to populate the the metadata for um trillions and trillions in tokens. It would take about 9 to 11 days, which is unacceptable, right? You if you're running this, you're probably running on millions of dollars in compute. You cannot have them um just waiting for you to fetch metadata for for 11 days.
11:25 So, the the um way that we've went past this is that um instead of of kind of um kind of going to S3, you can actually do um requests that are um batched together. So, instead of doing one at a time, you can batch them all together. S3 supports this through their APIs. And then this way um you can you can do this uh at at a massive scale by kind of increasing the the um the page size of your if your batches to to kind of the limit that that S3 gives you um which is about a thousand um kind of requests in one list and
12:02 that gets you to about um kind of two hours of API processing time to get your your metadata. So, first lesson, um, make sure that you know how you manage your metadata and make sure that you have, um, kind of thought through just getting the state of your S3 buckets into whatever synthetic data system that you have in order for you to, um, kind of be able to to even start doing this at scale.
12:24 So, 11 days to two hours. Uh, that's that's the first one. Um, awesome. So, the second one is GPU instability. So whenever you're running um these types of large scale synthetic data runs, of course you're going to hit some sort of bad GPUs, right? So that means that the machine might go down. Uh GPU might be bad. Something around those lines. Everything that you can expect when you try to run it at scale would happen.
12:50 So that leads to loss progress that leads to kind of wasted compute, partial outputs, um if you if you only parse through half of the file. And of course at the end of the day you're going to uh need more cluster time which is always hard to get. Um so the the the way that we do this in in our pipeline is that we have a Spark job that outputs the the the prompts together with the the documents that we need to rephrase that it's it stitches them together.
13:14 Um and then we start a rig job that that can actually uses the GPUs, right? Um so the the problem is that there's some sort of onetoone mapping between input and output. That's that's kind of how we we usually structure these types of jobs. Um, in practice, if your partitions are um quite large, what happens is that um you're going to do a bunch of of of um kind of progress.
13:44 And then if you if if a refrazing job on one node for one file for one partition takes eight hours and the box goes down in the last five minutes of those eight hours, you're going to lose a lot of time. So you need to think about um kind of you either have to checkpoint and whenever you do progress within the synthetic data um jobs you need to periodically say okay let me stop and like flash this back into S3 or um you need to have partitions that are small enough that you can tolerate that downtime.
14:14 So um you need to to think about either your partition size or in general kind of have some sort of checkpointing. Um what we've done is is a bit of both where we've managed to right size our partition sizes and checkpoint as well when it comes to this things. So um basically whenever um a job fails it gets retrieded but the amount of um kind of time that we lose is it's minimal and it makes sense for the setup that we have.
14:38 Um so from multi-day job failures now we have recoverable retries and it it it works it works as as expected. So just to recap resumeum execution some sort of IDM potency where whenever you're checkpointing you can override the file that you're checkpointing into and then you need to have some sort of mechanism for for simless recovery um cross infrastructure orchestration.
14:59 So in in our case the the the cluster where we origin run curation in general our pipeline runs in a centralized location and then customers might have synthetic data um jobs that run in various GPU locations throughout the world right you might it's going to be hard to get especially for synthetic data where you don't need necessarily fast networking in between your GPU nodes to get a lot of capacity in one place so what happens then is that you need to be able to orchestrate all of these jobs um in in various
15:27 clusters so it's not just one GPU cluster that you're working with you're working with um kind of a number of them. Um the problem here is that whenever you have um and you're trying to to to schedule um these jobs is that you might be able to um kind of schedule your CPU requirements for your ray jobs but not necessarily your um kind of GPU requirements in the GPU pool.
15:49 So you need to think about kind of being able to schedule in some sort of fashion across these two things. Now um once you have the the next job coming in um you have GPU pull capacity but you don't have CPU pull capacity right so in general whenever you're trying to do this this type of cross infra or orchestration you need to think about scheduling and you need to think about um kind of being able to schedule this effectively um and in this case as I mentioned you you you will get to a point where you cannot schedule
16:20 effectively so what happens then is that um you you kind of need to think about the the ways in which you do this in a in a proper fashion. So some sort of bin packing problem at the at the end of the day. Uh so how did we solve this? Um so same setup as before. Um we have this orchestrator that runs in a different cluster. We have the compute cluster that runs spark and the the ray jobs that mostly need CPU and then we have the the one that has H100s.
16:44 Um so whenever we schedule um we have dedicated pools for the specific resources that you need in those clusters. So if you have a spark requirements then you spark you might have a a spark pool that that has CPUs and then you dedicate special pools for um kind of the the ray head. So when you think about um some sort of atomic scheduling you need to think about both your CPU and your GPU requirements in order to do this optimally.
17:11 So in this case whenever we have more CPU um requirements that we need to schedule because we've decoupled these two things um and and then we can schedule the the um head and the worker for array um in a in a more atomic fashion then we don't have the problems as as before um and then everything scales as expected. So cross in front orchestration um you need to think about multiple providers multiple um clouds that you run this in you need to think about scheduling and u make sure that whenever you schedule CPU and
17:44 GPU resources especially for the same type of job you do that in a way that that makes sense. Um last up is inference optimizations. Um so whenever you run VLM or whenever you run SG lang um you need to think about the the hyperparameters that you use for your jobs right and uh what what we found is that you need to be u quite disciplined about this and you need to be able to to benchmark this and optimize and have a harness that just focuses on the VLM parameters you can get um a lot of gains um in the case in this
18:13 case about 40% throughput gains just by tweaking the flats for VLM um and SG line right so in our case we we went with VLM Right now we're actually testing HGLang as well. We do a lot of batch inference. Um which is very much different from from the the normal online one that everybody is used to. Um but whenever you you kind of orchestrate this through VLM, you need to make sure that you have the the right batch size, you have the right speculative decoding um setup, etc.
18:40 Um if you do kind of a bunch of grid searches over the parameters that you need to optimize, turns out you can get quite a lot of uh of results in terms of throughput. So to recap, um whenever you you're trying to do things, make sure that um you think about kind of uh the scale that you're running at. You have a good system for auto recovery whenever jobs go down.
19:01 You need to orchestrate in some sort of atomic fashion across CPU and GPU resources. Um and you need to make sure that that you kind of sweep over your VLM parameters. So what what this setup helped us achieve is is quite straightforward. Um and in at the end of 2014, sorry, 2024, uh we can only output around 30 billion synthetic tokens in kind of limited runs within our SERM cluster.
19:25 Uh our requirements for synthetic data went up by by a bunch. And then um kind of end of last year, we were able to to run this massive uh synthetic data jobs that that um generate in this case for web about 7 trillion tokens and then for math and code another five for a bunch of our customers. uh we're continuing to ablate and and kind of search for better synthetic G recipes.
19:45 Um we're actively working towards this. So keep an eye out. We'll be publishing more. And um I think domain coverage is quite important. I think something that we haven't talked today is how does this differ in between web, multilingual, math code, uh legal data, etc. By the way, we're hiring um throughout the board for data infrastructure, cloud infrastructure, normal infrastructure, program managers, etc.
20:06 So hit me up afterwards. Uh thank you for the team that worked on this and uh thanks everybody for attending my talk. [music]