tldr;
This summer at Eventual, I helped build the data pipelines behind a scenario-mining platform for one of the largest Physical AI companies in the world. This article talks about how we went from a manually run pilot, to an automated orchestrator I designed and built myself, to a production design built on Apache Iceberg, with versioned tables, change data capture, and incremental cross-table ASOF joins.
Intro
During my summer at Eventual, we built a scenario-mining platform in collaboration with one of the largest Physical AI companies in the world. The goal was to enable researchers & safety leads to find the right data to post-train their models on. The platform would let you search over your entire corpus of data for specific scenarios, think “robot presses the square peg into its hole and rotates it back and forth to find the alignment before it seats”, and find the exact moments where it occurred in your data.
A core piece of infrastructure supporting this platform was our data pipelines, which involved ingesting a customer's data, computing features on it (e.g. embeddings, object detection) before serving a downsampled version on our platform. Since robot data like camera frames and telemetry are recorded at different rates into their own tables, that last step involved an ASOF join* across tables, which comes with its own complexity (more on that later in the article).
*ASOF joins were coincidentally also a feature I previously built in Daft, our high performance data engine! Read how I did it here.
We wanted this pipeline to be done automatically, incrementally, at scale.
By the end of my internship I had helped build out three versions of this. One for the pilot, one to tide us over while we transitioned to production, and finally one that we would confidently deploy in production. In this article, I’d like to walk you through all 3.
V1: Why was this hard?
The first version of the pipeline we built out during a pilot (in a month!) had a few issues*
- There was no dependency tracking. Each job assumed its parent had already finished. If the parent was still running, the child would silently compute over incomplete inputs. So instead of running jobs concurrently, the intern (me) had to watch for each one to finish before kicking off the next.
- Debugging was a mess. We had yet to integrate any real-time status updates to our data jobs, which meant that the intern (me) had to babysit a run whenever it kicked off.
*This might seem like an overly naive design for a team of really smart data systems guys! The truth is, a pilot is an extremely short window for a company to prove itself, and different engineering goals have different engineering priorities. Tech-debt always needs to be strategically chosen in order for a product to succeed, and I felt like we actually did this well!
V2: What did I do about it?
The short-term goal was to turn this manual process into an automated one, in the simplest way possible while still maintaining correctness. Here’s how:
The first order of business was making sure all of our jobs could run concurrently with each other. At plan time, each job only selects partitions whose parent inputs are already complete, and anything not ready yet is simply picked up by a later run.
We could then set up a cron function for each job that would run at a fixed interval. I wanted the cron function to act as a state machine: each time it fired, the function would check a durable dictionary to see what state the job was in. From there, the function did one of three things:
- No run in progress: plan the next batch of ready work, spawn workers, and record the run + its workers in the dictionary
- Run in progress, workers still going: poll each worker's status and post a progress update.
- Run in progress, all workers reached a terminal state (i.e. success / failure): post a final summary, and clear the dictionary so the next invocation can plan a new run.
A bunch of things worked in our favour: our data jobs were idempotent, which meant that we could always retry safely without needing to guarantee exactly-once dispatch. I also added two quality-of-life improvements:
- notifications were sent to a slack channel with links to the logs of failed jobs (if any)
- jobs that repeatedly failed were marked as blocked in our dictionary, preventing us from wasting compute until a human manually intervenes.
We now had an automated data pipeline!
V3: Production
This was a big improvement, but it still had two flaws, both revolving around the ASOF joins in the last downsampling step mentioned at the start:
- Joins couldn't be done incrementally: Computing features incrementally was easy, since we just filled in whatever rows were missing. But a joined row depends on rows in several tables, and those tables were constantly being backfilled, deleted and upserted. Figuring out which joined rows a change affected meant scanning everything, which was not feasible at the scale of data we were working with. We ended up re-running the entire join whenever enough had changed.
- There was no versioning: With jobs running concurrently, a join could read fresh data from one table and stale data from another.
For production, I built on the design work of two extremely talented senior / staff engineers, Chris & Rohit. At its core was a simple idea that fixed both flaws: every batch of ingested data was assigned a version number, which every table downstream of it inherited.
This unlocked two things:
- Consistent views: A join no longer reads whatever happens to be in each table right now. It asks for all of its inputs as of a single version. Every table carries the same version numbers, so the join is guaranteed to see the same batch of ingested data everywhere. If a new batch lands halfway through, the join won't see it until its next run.
- Incremental joins: Since every version is kept, we can ask "what changed between version 41 and 42?" and get back exactly which rows were added or deleted, including what the deleted rows contained. The join then only recomputes the rows affected by those changes, instead of scanning everything. Each job just remembers the last version it processed and picks up from there.
For this to work, versions needed to be cheap to create (no more data duplication please!), read, and compare. We opted to use Iceberg as our storage backend, as that's more or less what it was built for.
A short TL;DR on Iceberg
Iceberg splits a table into two layers:
- Storage: the actual data, stored as immutable Parquet files in object storage (e.g. S3).
- Snapshot metadata: a small set of files, each describing a snapshot of the table at a point in time and the data files that make it up. In our case, each snapshot is a table version.
For us, this meant that:
- Creating a version is cheap: Data files are never modified, so a new snapshot just points to the newly added files plus the unchanged ones. Nothing gets copied.
- Reading a version is cheap: Each table records which version each of its snapshots corresponds to, so reading "as of version 37" resolves to the right snapshot in every table and reads from it.
- Comparing two versions starts with their metadata: which files were added and which were removed. We only need to open those files, not the whole table.
From a performance perspective, separating storage and metadata helps with:
- query optimizations, like using metadata to prune the search space.
- scaling writes, since workers can do the heavy work of writing Parquet files in parallel, while a single committer handles the small metadata commit, publishing files from many workers at once.
When it comes to correctness, we needed to make sure that a table could never end up in a bad state, and three different properties of the system got us there:
- Atomic commits: A new version only becomes visible once its snapshot metadata is swapped in, and that swap either fully happens or doesn't happen at all. Readers never see a half-finished write. If a commit fails partway, the table stays at its previous version, as if nothing happened.
- Idempotent commits: Every commit states which version it expects to build on, along with an idempotency key. If a worker retries a commit that already landed, we recognise the key and skip it instead of writing the data twice. If a writer is working from a stale version, its commit is rejected.
- Recovering from partial commits: A single ingest usually writes to several tables, and sometimes only some of them commit before something fails. Every table records which version it's at, so the ingestor can check which tables fell behind and resume just those, instead of starting over.
Reflection
Looking back, none of these versions were wrong. Each one was the right design for what we needed at the time. The pilot had to prove the product in a month, so we took on the tech debt we could afford. The stopgap had to take the humans out of the loop without slowing the team down. And production had to be correct at a scale where re-running everything was no longer an option.
Huge thanks to Chris, Rohit, and each and every one of the engineers at Eventual for an incredible summer.