Francesco's avatarFrancesco Ceccon

A new Iceberg Storage Engine

I spent the past year researching storage engines built on object storage, I believe they're the future to store agentic data. We have already seen agents pushing existing infrastructure to its limits, the most familiar example for developers is GitHub. Their usage increased by orders of magnitude, that took their existing services down because they were built for human-level usage.

For a company, the first step to understand how agents use their product is to collect data. This is nothing new. What's different now is that startups need to do it earlier and at larger scale than before to stay competitive.

The current data ecosystem is built for larger startups and enterprises. Tools are complex, with many moving parts, and infinitely configurable.

What's missing is an engine that does what needs to be done and nothing more. Minimal compute overhead. An engine that scales when the catalog has millions of tables. An engine built on modern primitives, that is easy to operate after one year as in day one.

Does it have to be this way?

The first question I asked myself when embarking on this journey is if these tools are complex because of their history or by necessity. I believe it's the former.

Take how data first enters a pipeline (ingestion). What we need is a place to make data durable so that it can be later materialized. Ideally, we would want this system to also validate incoming data, but this is not always possible. In most deployments today this is done through Kafka. Kafka is an amazing piece of software but it's also overkill if all you need is a temporary place to park data before it goes into a table. Companies realized this is silly and we're seeing products like DataBricks Zero Ingest hit the market.

We know now it's possible to build a write-ahead-log (WAL) completely on object storage (e.g. see this repository by Jack Vanlightly for a list of implementations). So what if we did just that and had another process read from the WAL to materialize the table? We would get a system that's simple to operate (because it's stateless on object storage) and that can scale very well (again, because it's on object storage). This is exactly what we want if we have many autonomous systems producing data constantly. We want to experiment new agentic workflows freely without worrying if our infrastructure can handle it or if we're getting a surprised bill from our data platform. We want to collect data from any tiny subsystem in our company so we can feed it to the latest and greatest AI.

The good news is that I believe it's possible to build such a storage engine for Apache Iceberg. Iceberg is the target because we want data to be accessible by any tool. The feature set is minimal, we only implement what's necessary:

  • Data is validate before ingestion. Clients get an error immediately if they push bad data, if it's an agent it means they can fix it and retry.
  • Rows are deduplicated by primary key. There is no need to pull in a streaming engine just for this.

With just these two features we can easily spawn tables that only have valid, deduplicated rows (if you're familiar with the Medallion Architecture, it means we store Silver, in some cases Gold, data in one step).

The design

The design is heavily influenced by projects like Apache Paimon, Apache Hudi, but modified to use Iceberg V3 primitives.

As I mentioned before, the ingestor validates incoming data and periodically flushes it to a WAL on object storage. Once data has been flushed, it's durable. For this, we use the Arrow Flight protocol since it's efficient and widely deployed already.

Table's data is stored in one or more base files, data in base files is deduplicated by primary key. Rewriting these files every time new data is inserted would be too expensive, so materialization creates "delta files". Delta files are a combination of Parquet (data files) and Puffin (deletions, including replaced values) files. To improve read performance, the engine periodically compacts these files into new base files. Since the engine has a completed view of the files in the system, it performs all the table maintenance you expect (snapshot and data pruning) and more (splitting base files once they grow too big).

We're building on top of the existing Rust Data ecosystem, so we leverage the thousands of hours of work contributors put into Arrow, Parquet, DataFusion and Iceberg.

I'm working on this

I started working on this design and it's now on GitHub, you can follow along as I work on it.

I will explain the core architecture in a later post, but the idea is to have a reconciliation loop that drives table's state forward until all tables have the most recent data and optimal file layout published to the Iceberg catalog.

Thanks to Thomas Hsueh for reviewing drafts of this.