We recently gave a talk at the NYC Iceberg community meetup about our work on physical replication of Iceberg tables. Here’s a summary of the presentation.
Today in the world of data lakes and OLAP, open table formats (like Iceberg) are a dominant theme at data warehouse conferences and on product roadmaps. As the name implies, open table formats (OTFs) unlock a level of interoperability at the storage layer far superior to native proprietary storage formats all without sacrificing performance. But there are other, less obvious benefits of systems built on OTFs, and I think our use case is a great example of one of them.
At Prequel, we work with software companies to help them share data with their customers. We do this by powering “connections” between data platforms, and we do our best to use the optimal method for each pairing, given a number of constraints. As adoption and support for Iceberg grows, we get access to new opportunities for the design of these connections, one of which has been particularly exciting: table-aware, partition-aware, physical Iceberg replication. Early results from this new design have shown throughput improvements of close to 10x, and we want to share a little about how we did it.
Data federation vs. data replication
First, there are two related concepts to distinguish: data federation (“zero copy”) vs. data replication. Assuming the data producer and data consumer reside in at least two clouds or regions (a safe assumption for software businesses and their customers – our use case), we must weigh the tradeoff between the two. Even if we assume every data platform will have perfectly interoperable Iceberg support, the cost of data egress is one of the governing considerations in the world of data sharing. Federate across regions: operationally simple, but incur incremental egress on each table read. Replicate across regions: additional upfront setup, but egress costs are predictable and capped.
Strategies for replication
At the highest level, there are two broad categories of data replication: logical and physical. Logical replication works with data at the semantic level, with no regard for how the data is stored on disk. This usually means continuously polling the source query engine for changes to the data, and replaying those changes to the query engine of the target system. If the engines use different storage formats, this might be the only way to translate tables between systems. Alternatively, physical replication involves moving the underlying bytes between regions as they appear on disk, which can avoid the expensive and inefficient work of reprocessing every data file.
For Prequel, we started with logical replication, and it worked pretty well for a while. Query engines on most mature OLAP systems provide a number of features to make unloading/loading data pretty efficient, even without access to the underlying storage, and we’ve taken advantage of most of the obvious opportunities for optimization. As an example, we have individual table jobs that reliably replicate over 50 billion rows (~4 TB) every day using logical replication. Still, some of our customers asked us for more throughout, and after all of the obvious optimizations were implemented in our logical replication engine, we started looking into opportunities to implement physical replication.
Physical replication is only possible if a replication job connects two data platforms that use the same storage format. Prior to widespread OTF adoption, physical replication was a feature reserved for proprietary engines writing to their managed storage and only made sense within the same query engine (those engines alone had access to and understood the physical storage contract). But today, open formats like Iceberg create an opportunity for specialized cross-cloud, cross-region replication engines.
Table-aware physical replication
Open table formats are, at the core, detailed specifications for correctly writing data into interoperable storage. The implementation, however, is left to the engines that adopt them. For Iceberg, the specification details how to structure the data and metadata files in object storage so that query engines can read the tables. Because the data and metadata are fully contained in object storage, it may be tempting to think of Iceberg replication as a bucket-level replication job that moves both between regions. However, the Iceberg spec has strict referential integrity requirements (and many object references), so naive bucket-level replication would fail for any real world workloads.
Table-aware replication is the “correct” way to build Iceberg physical replication. Table-aware replication involves an understanding of the Iceberg spec in order to read from the source table and generate a derived, spec-compliant target table (ideally without rewriting any of the data files). Replicating a full table from one region to another can be as simple as iterating through metadata, enumerating manifests, manifest lists, and data files, moving everything in order with minor metadata rewrites, and registering the table in the target catalog.
But Iceberg is a complex spec, and anything more advanced can get a little more involved. For our use case, things did get a little more advanced.
Partition-aware replication
Basic table-aware replication assumes every row of data in the source table needs to be replicated to the target. But as discussed, moving data more than absolutely necessary is expensive, and for our use case, each row in the source table usually only belongs to one target cloud and region. In an ideal world, each row of data is copied at most once.
Because Iceberg supports physical partitioning, a source table with partitions configured should mean that Iceberg engines are automatically grouping those partition predicates into the same data files.
For our use case, the most important predicate is the “tenant” (the owner of a given row of data). If each tenant’s data is physically partitioned, in the best case scenario, physical replication of data files should still work without reprocessing. Metadata files aren’t as cleanly segregated, but they can be rewritten more cheaply and with minimal data processing.
Implemented correctly, Iceberg tables with thoughtfully designed physical partitions should allow for efficient, table-aware, partition-aware physical replication. Implemented incorrectly, and this becomes more than an efficiency concern: it creates a data privacy risk. If Customer A’s data shows up in Customer B’s data lake, the consequences will be far greater than the incremental cost of egress.
Things to look out for
By reading from and writing directly to Iceberg storage, you not only take on the responsibility for interpreting and following the Iceberg spec, you also take responsibility for the appropriate governance of that underlying data. This means translating user intent into the relevant storage operations to produce and maintain a correct Iceberg table replica. Concepts like snapshotting, schema evolution, partition evolution, row and column deletion must be deeply understood and correctly implemented.
Iceberg is a fully featured spec, and has a lot of advanced capabilities that make it popular and user friendly. The flipside is that reading and writing while adhering to that spec and all of it’s features, can be difficult.
For our use case, there were a few areas we had to look out for:
- Versions. The Iceberg spec is evolving, and techniques that work in V2 and V3 may need to be adapted for V4+.
- Partitions. As we already discussed, physical partitions are particularly important for our use case. But Iceberg supports partition evolution, so physical partitioning cannot be verified just once and must be validated for each manifest.
- Snapshots. Snapshotting makes table history a convenient built-in feature of Iceberg tables, but it can complicate the data replication story. Suppose Provider 1 wants to replicate data to Company A: would the provider expect the full table history to be included in the initial replication job?
- Deletions. As with most OTFs, row deletions are represented as “soft deletes” or “tombstones.” Query engines should respect them, but if the data is replicated outside the governance layer, soft deletes may unintentionally expose deleted data.
- Schema evolution. Similarly, schemas can change over time. The catalog should hide inactive columns, but those columns may remain unused in the underlying data files.
How we did it
As a guiding principle, in our role as an external engine with our own governance responsibilities, we designed an intentionally conservative approach for table-aware, partition-aware Iceberg replication.
In this case, “conservative” means “failing closed” for any Iceberg table state that could lead our users to move data unintentionally or nondeterministically. That means inspecting every physical object in flight and rewriting it when necessary to ensure an accurate and safe replicated Iceberg table.
For each layer of an Iceberg table in object storage, we evaluate every line of metadata to determine which to keep, rewrite, or delete entirely (based on the categories we highlighted above).
The results
For some of our highest-throughput workloads, physical Iceberg replication improved table-level throughput by 9x, moving 10,000 data files across clouds in under three hours. This result was more than enough to meet demand from some of our highest-volume customers – all without additional tuning or workload-specific optimizations. We’ve already identified several such optimizations and are working on them. This should unlock the next tier of service our customers have been asking for and open the door to further improvements that increase throughput even more.
What’s next
Though we’ve supported open table formats like Iceberg and Delta for a while now, the critical mass of platform support is relatively recent, and this work is evidence of some of the value we can deliver as adoption continues to grow. We’re excited to continue improving our tech alongside the other ecosystem partners that are improving their own support. If this work is interesting to you, we’re hiring engineers to help us accelerate it. Reach out to careers @ and let’s chat.