Weather Forecast Data: Icechunk vs. Iceberg Head-to-Head
CEO & Co-founder
TLDR: Operational weather forecast data workloads built on Icechunk are 10-100x faster and 5-8x cheaper than Iceberg.
Warning: This article is long. It’s written for a system architect planning to expend significant budget on weather data infrastructure and looking to deeply understand the options and tradeoffs. It is based on an empirical study whose design is documented and openly published for the purpose of independent validation.
Introduction
Weather forecasting delivers critical information for agriculture, energy and utilities, aviation, transportation, maritime shipping, emergency management, insurance, construction, retail—not to mention the importance for defense and national security. The impact is hard to quantify globally, since weather forecasts are so embedded into our daily lives, but London Economics estimated that the UK Met Office will deliver £56 billion of economic value to the UK alone over the next decade.
Today, weather forecasting is undergoing a rapid transformation, as new players and technologies enter the game. Tech companies both large (e.g. DeepMind, NVIDIA, Microsoft) and small (e.g. Brightband, WindBorne, Jua, Zeus AI, Planette), in addition to innovative government agencies like ECMWF, NOAA, and the UK Met Office, are deploying AI powered forecasts that are faster, cheaper, and increasingly more skillful than conventional numerical weather prediction models. This activity is producing a deluge of weather forecast data.
Furthermore, advances in cloud and database technology, together with agentic coding, are making the raw forecast data more accessible to non-specialists.
Taken together, this means that more and more companies are ingesting and storing more and more weather forecast data than ever before. They are comparing, scoring, and evaluating different models for their specific business needs. Many are even fine-tuning or outright building their own custom AI weather models. They are discovering what the big agencies already knew: weather data are absolutely massive (ECMWF alone generates 400 TB of data per day), and building efficient data systems to turn the raw data into business value is no picnic.
Note: This analysis was done with weather data and example queries used by data scientists working with weather data, but it is broadly applicable to other multimodal, n-dimensional array datasets, such as observations that originate from sensors and satellites, microscopes and imaging systems, and scientific simulations more generally.
At Earthmover, we work with these types of organizations every day, from innovative AI startups like Brightband to large government agencies like the US National Weather Service.
The point of the analysis shared in this blog post is to answer a question we encounter frequently, especially from experienced data engineers who are unfamiliar with weather forecast data: should I store weather forecast data in a tabular data format, specifically Apache Iceberg?
We see many organizations go down this path because it’s a safe choice and then pay heavily for it. As shown by the analysis below, Icechunk offers a much lower total cost of ownership, plus far superior performance, for weather forecast data systems.
Objectivity via Agents
We’re obviously and transparently biased because we developed Icechunk. But we’ve also done our best to be objective. Leveraging agentic coding provided a simple way to level the playing field. We defined the basic parameters and questions for the study, but the implementation details were left to Claude. (This is how most data engineering teams are working today.) The choices made, and the resulting outcomes, reflect the level of expert knowledge about Iceberg, Icechunk, DuckDB, and Xarray that are baked into Fable 5.1 and Opus 5.5 at medium effort—not some esoteric knowledge that is unique to our team. Of course, we carefully reviewed these choices, and we performed several passes of iteration to make sure each option was given its best possible shot before publication.
To conduct this study, we took one week’s worth of real operational ensemble forecasts from the UK Met Office, stored it six different ways, and measured the total cost of ownership (TCO): what it cost to ingest, store, and query the data under simulated but realistic access patterns. We ran everything in AWS, and the code to reproduce the results end-to-end is open source, along with the prompts and harness we used to develop it.
%%{init: {"flowchart": {"wrappingWidth": 260}}}%%
flowchart TB
src["<b>Source data</b><br/>UK Met Office MOGREPS-G<br/>1 week, 28 forecast cycles<br/>6.5 TB of NetCDF on S3"]
m_dl["Download<br/>the files"]
m_fuse["FUSE mount<br/>the files"]
m_ib["Transform to<br/>Iceberg"]
m_vz["Virtual Zarr /<br/>Icechunk"]
m_nz["Native Zarr /<br/>Icechunk (2 layouts)"]
queries["<b>Three queries</b><br/>Q1: point timeseries<br/>Q2: regional statistics<br/>Q3: ML dataloader"]
metrics["<b>Metrics</b><br/>Storage cost<br/>Ingestion cost and latency<br/>Query cost and performance"]
tco["<b>Total cost of ownership</b>"]
src --> m_dl & m_fuse & m_ib & m_vz & m_nz
m_dl & m_fuse & m_ib & m_vz & m_nz --> queries
queries --> metrics --> tco
Background: Tensors vs. Tables, Iceberg, and Icechunk
A critical thing to understand about weather forecasts is that the data are fundamentally tensor-shaped.
For each predicted variable, a global forecast model produces a value at every grid cell on the globe, for every forecast lead time, for every initialization time.
So the temperature output of an ensemble forecast is naturally a five-dimensional array: (init_time, lead_time, member, latitude, longitude).
There are dozens of such arrays, one per variable, and they all share the same coordinates.
We wrote about the consequences of this structure last year in Fundamentals: Tensors vs. Tables. In that post we showed this with a small, single-forecast example (4 GB) and found that Xarray + Zarr beat DuckDB + Parquet by up to 10x on selection queries. But a 4 GB benchmark on one forecast does not tell you much about running a production system on tens of terabytes. This post repeats a similar comparison at a much larger scale, with a proper accounting of cost, not just speed.
Two storage technologies are central to the comparison.
Apache Iceberg is an open table format for data lakes. Iceberg turns a collection of Parquet files in object storage into something that more resembles an actual database by enforcing consistent schemas for tables, organizing partitions, and, most crucially, providing ACID transactions so multiple concurrent processes can safely read from and write to the tables. Iceberg has become the default answer to “how do I store large structured data in object storage?” Every major tabular query engine and data warehouse can read it.
Icechunk is an open storage engine for Zarr, the leading cloud-native format for chunked multidimensional arrays. It borrows the ideas that made Iceberg successful—snapshots, transactions, a manifest layer between the client and the object store—and applies them to tensors rather than tables. Icechunk can also store virtual chunks: references to byte ranges inside existing files (NetCDF, HDF5, GRIB) rather than copies of the data.
Both formats are structurally similar, requiring only object storage as infrastructure. (Actually, that’s not quite true—Iceberg requires an external catalog service to coordinate and serialize transactions, while Icechunk uses the object store’s built-in conditional write primitives, but we’ll ignore that detail for now.) They differ in the data model; Iceberg is tabular, while Icechunk is tensor-based.
Because the data originate as NetCDF files, we also benchmark approaches that bypass the table formats entirely and just download the files directly as needed to satisfy the queries.
The Rules
Our experiments were designed to respect the core cloud-native principles that have made both formats popular in the first place:
- Object storage as the only persistence layer. On the Pareto frontier of cost and performance, object storage has won the day. It’s 10x cheaper than volume-based storage while providing astonishing horizontal scalability and architecture flexibility. The downsides, like higher time-to-first-byte latency, have largely been mitigated by cloud-native formats. If you’re storing Terabytes of data or more in the cloud, object storage is really the only reasonable choice.
- Multi-engine support. A big advantage of open table formats is to avoid lock in to a single vendor or software to read your data. These formats can be queried from a wide range of different engines, libraries, and languages. This is an important ingredient to data sovereignty, an increasingly critical topic for government agencies.
- Just run queries from a EC2 node. There are many sophisticated distributed query engines out there, from Snowflake and Databricks to Earthmover’s new Zax tensor query engine. However, under the hood of all of these is basically an EC2 node talking to S3. A great advantage of the open table formats is that you can always just fire up an EC2 node next to your data, open your preferred open-source analysis tool, and get decent performance off the bat. Conversely, if a storage format can’t do well in this simple scenario, scaling out to complex multi-node architectures is probably a waste of money.
So our experiments are simple: run all of the ingestion and queries from raw EC2 nodes using standard open-source libraries. We chose Xarray and DuckDB as our tools of choice—each represents the industry standard tool for ND-array and tabular data respectively. (The results are not highly sensitive to this choice, but we’ve got a more thorough comparison of engines planned for a future post.) Furthermore, we disabled local file caching for all our engines, which can strongly bias results when queries are repeated within a benchmark suite.
The only two AWS services which accrue costs in these experiments are S3 and EC2, making it easy to translate storage and performance metrics directly to dollars.
The Original Data
We used MOGREPS-G, the UK Met Office’s global ensemble forecast, which the Met Office publishes on the AWS Open Data registry. It is a good test case for a few reasons. It is a serious operational product from a major agency. It is big but not absurdly big. It is delivered in a format (NetCDF) and a layout (one file per variable per lead time) that is completely typical of how weather agencies distribute data. And it is directly relevant to many of our customers.
MOGREPS-G runs four times a day (00, 06, 12, and 18 UTC).
Each run, which we call a cycle, produces 18 ensemble members on a 960 x 1280 global grid (about 20 km resolution), at hourly lead times out to 132 hours and then 3-hourly out to 246 hours, for a total of 171 lead times.
The Met Office publishes each variable at each lead time as a separate NetCDF file containing all 18 members.
A single file is about 35 MB compressed (88 MB of float32 values), stored internally as HDF5 chunks of shape (1 member, 128, 128) with zlib compression.
A single cycle of surface variables is 11,521 files and 232 GB.
The bucket keeps a rolling archive of about a month.
For this study we took one week of forecasts: 28 cycles from 2026-09-12 to 2026-09-18, restricted to the 81 surface-level variables (we left out the 3-D pressure-level and height-level fields). That is 322,588 files, 6.5 TB of compressed NetCDF, and about 28 TB of uncompressed float32 data. Flattened to a table, it is 106 billion rows.
Because the source bucket is a rolling archive, we first copied our week to a bucket of our own in us-east-1, so the results remain reproducible after the Met Office deletes the originals.
Every method below was built from this static copy, and all compute ran in the same region.
A continuous operational implementation would cycle data in and out of the Iceberg / Icechunk stores dynamically, but that wouldn’t change any of the cost or performance metrics reported here.
The Queries
By “query” we mean a specific data analysis task which requires access to the forecast data.
Different storage layouts will perform differently for different queries, and the total space of possibilities is quite overwhelming.
Nevertheless, after observing usage patterns across hundreds of different organizations, we’ve boiled things down to three simple queries that are broadly representative of common workloads on forecast data.
Q1: Point Forecast Timeseries
This query represents an extremely simple but ubiquitous access pattern.
Extract the 2 m temperature forecast at a single location (we used London) for one cycle: all 171 lead times and all 18 members, 3,078 values. Then compute heating degree days for each forecast day from the daily mean temperature.
An energy trader, a retailer, or a weather app wants the forecast at their location, and they want it now. It reads a tiny amount of data (12 KB) out of a very large dataset, and it happens constantly. In our cost model we assume 10,000 of these per day.
This query also frequently appears on the back-end of systems which serve forecast data to clients via HTTP APIs, so it’s an important baseline for system architects.
Q2: Regional Ensemble Statistics
For a box covering the UK (54 x 37 grid cells), for one cycle, compute the ensemble mean, ensemble spread, and CRPS at every lead time for two variables (2 m temperature and 10 m wind speed).
Since we have no observations, we use the next cycle’s ensemble mean at the same valid time as a proxy analysis, so the query reads two cycles.
This is the forecast verification pattern. Every organization that consumes forecasts from several models eventually wants to score them against each other over their region of interest. It touches a medium amount of data (about 150 MB of values spread across two cycles) and involves real computation. We assume it runs once per cycle, four times a day.
Q3: ML Dataloader
This pattern is designed to simulate a typical workload in AI/ML research: feeding a training loop for an AI weather model.
For a randomly shuffled sequence of
(cycle, lead_time)samples, load four global variables (temperature, wind speed, pressure, humidity) for all 18 members as one float32 array, in batches of 8 samples. Each sample is 354 MB decoded, and the benchmark streams 16 samples (5.7 GB) as fast as possible into host memory.
This is modelled on the way frameworks like ECMWF’s Anemoi consume forecast archives. Unlike Q1 and Q2 it wants whole fields, not slices, and it is throughput-bound: the only metric that matters is how fast the GPU gets fed. We assume one training epoch over the week per day.
Methods Evaluated
We evaluated six different ways of storing the data:
File Downloads
The simplest (some would say naive) way to interact with this data is to just download the NetCDF files to a local disk before opening them. This method is the least cloud-native; it treats S3 as if it were a big FTP server. This is an important reference point, because that’s basically the assumption that most weather forecast agencies make about how users will interact with their data. While the cloud-native technology world has long since moved past this way of working, it’s still probably the most prevalent approach across the industry.
The obvious performance limitation on file downloads is that we can only download a whole file. For Q1, that means downloading 171 files (6 GB) to read 3,078 values (12 KB).
We implemented this as fairly as we could: list the cycle, download every needed file in parallel with obstore (we also tried the AWS CRT transfer manager; the two were within a few percent), open them with h5netcdf, compute, delete. There is no ETL and nothing to store beyond the NetCDF files themselves.
FUSE Mount the Files
FUSE (Filesystem in User Space) is a software interface that allows users to create custom filesystem implementations. A common approach when porting legacy applications to the cloud is to use FUSE to make S3 look like a local filesystem. That way, application code written for files doesn’t have to be rewritten.
The performance of FUSE can be highly variable and hard to predict; implementations each employ different strategies for read-ahead, caching, and parallelism which are mostly invisible to the user.
We used Mountpoint for Amazon S3, AWS’s own FUSE client, mounted read-only with no data cache, and read the files in place with h5netcdf. In principle this should be better than downloading, since HDF5 can seek to just the chunks a query needs. In practice, as we’ll see, it wasn’t.
Transform to Apache Iceberg
This is the safe choice made by data engineering teams who are unfamiliar with scientific data formats.
We tried hard not to straw-man this option.
We built the table the way a competent data engineer would: one wide table with a row per (init_time, lead_time, member, lat, lon) and one float32 column per variable (81 value columns plus six key columns), partitioned by init_time and lead_time using Iceberg’s hidden partitioning, written as zstd-compressed Parquet with 128 MB row groups.
That gives 4,788 partitions and 19,152 data files.
We used the Iceberg REST catalog built into Arraylake and queried the table with DuckDB, widely considered to be the best embedded engine for this kind of work.
We then spent a second pass optimizing it, and this turned out to matter a lot.
Our first version sorted the rows within each partition by member, lat, lon.
That is the obvious order, and it is terrible for point queries: the London grid cell appears once in every member’s block, so a point query touches every row group in the partition, and Parquet’s min/max statistics cannot prune anything.
Q1 took two minutes.
Sorting the rows along a Hilbert curve over (lat, lon, member) (so the 18 ensemble members of each grid cell sit in consecutive rows), puts spatially adjacent cells in the same row group; with this tweak, Q1 dropped 5x.
Everything reported below uses the Hilbert-sorted table.
A brief aside on the “geospatial” angle… A reader with a GIS background might reasonably ask whether GeoParquet or a spatial index partitioning scheme would lead to better results for the table.
The answer is no.
The table’s main challenge is not finding the right rows; it is the number of objects and requests it takes to reach them.
Our Hilbert-sorted layout already delivers what a spatial index would: within each (init_time, lead_time) partition, the rows for nearby grid cells sit in the same row group, so Parquet’s min/max statistics on lat and lon prune the point and box queries to a handful of row groups per file.
That is exactly what GeoParquet’s bounding-box covering does, minus the overhead of a WKB geometry column that is redundant on a regular grid.
But a single forecast timeseries query spans 171 lead-time partitions, each one requiring a footer read and a few range requests, on top of planning a query over 19,000 data files.
Adding lat/lon buckets as partition columns makes this worse, not better: each bucket multiplies the file count, so the point query still touches 171 files (now smaller) and every query plans over more of them.
The regional query has the same shape, spread over two cycles, while the dataloader query reads whole partitions and is indifferent to spatial layout entirely. The fundamental issue is that a forecast has three or more orthogonal dimensions that queries slice independently, and a table can be physically ordered along only one path.
Chunked arrays do not have this constraint; a chunk shape is a partitioning along every dimension at once.
Virtual Zarr / Icechunk
The virtual Zarr approach keeps the original NetCDF files exactly where they are and builds an Icechunk repository containing only references: for every chunk of the logical 5-D array, the file, byte offset, and length where those bytes already live.
VirtualiZarr does the scanning and Icechunk stores the references.
The result looks and behaves like a single Zarr dataset with dimensions (init_time, lead_time, member, lat, lon) for every variable, and Xarray opens it in one line.
No data is copied.
The catch is that the chunk layout is inherited from the source files.
The Met Office wrote its HDF5 chunks as (1, 128, 128), about 40 KB each, and that is what a virtual Zarr reader gets.
A point timeseries needs one 40 KB chunk from each of 171 files x 18 members, about 3,000 tiny S3 requests.
Whether this is a problem depends entirely on how quickly the query engine and underlying hardware can process all those requests.
Virtualization obviously has an important cost advantage: you don’t have to copy the data. Storage is not a massive line item at 6.5 TB, but extrapolated up to larger archives, this can become a major win.
Native Zarr / Icechunk
Finally, the native Zarr approach: read the NetCDF, rechunk into shapes chosen for the workload, recompress, and write a new Icechunk repository. This is the tensor equivalent of the Iceberg transform. The price, as with Iceberg, is a full copy of the data; in return you are unbound from the chunking decisions made when the NetCDF files were written.
For our native Zarrs, we compressed everything with pcodec, a lossless numerical codec that compressed this data nearly 8x (vs. 4.4x for the Met Office’s zlib and 5.5x for the Iceberg table’s zstd). We also used Zarr sharding so that many small chunks live inside one large object and can be read with byte-range requests.
Choosing chunk shape is analogous to choosing partitioning, and there is no single shape that is best for everything. So we built two layouts.
- Timeseries layout (
native_ts), following the pattern dynamical.org uses for its GEFS archive: one array per variable, chunks of(18 members, 57 leads, 16 lat, 16 lon), about 1 MB each, 12 shards per variable per cycle. A point timeseries query reads just three chunks. - Dataloader layout (
native_dl), following the ECMWF Anemoi pattern: one array holding all 81 variables, with one 7.2 GB shard per(cycle, lead_time)sample containing 1,296 inner chunks of(18 members, 240 lat, 320 lon). A training sample is a handful of contiguous byte-range reads from one object.
Both layouts hold the complete dataset, so a team that wanted both would store the data twice. We report them separately and also discuss the cost of supporting both simultaneously.
The Metrics
We measured three things for every method: how much it costs to store, how much it costs to build, and how fast it answers the three queries. We combined all of these into a monthly bill.
Prices are AWS us-east-1 list prices in September 2026: S3 Standard at $0.023/GB-month, a c7i.16xlarge (64 vCPU) for ingestion at $2.856/hour, and an m7i.4xlarge (16 vCPU, 6.25 Gbps network) for queries at $0.81/hour.
All timings are medians over repeated runs on an otherwise idle instance, with every cache disabled so that every run reads from S3.
Storage Cost
A few things stand out.
The Iceberg table is smaller than the source NetCDF.
Flattening the coordinates into 106 billion repeated (lat, lon, member, lead) tuples sounds disastrous, but Parquet’s run-length and dictionary encoding compress them nearly to nothing, and zstd beats the Met Office’s zlib on the data variables.
The 17% of value cells that are NULL (because different variables have different lead schedules) cost nothing measurable.
The native Zarr repositories are the most compact, at about 56% of the source. That is 30% smaller than Iceberg for the same values. This is due to pcodec’s smart algorithm but also the fact that n-dimensional chunks preserve more data locality, resulting in more efficient compression. The compression ratio is important. Using state-of-the-art compression can effectively buy us a whole second copy of the data, which can then be used to support query patterns which don’t align well with the original chunking.
The virtual repository is essentially free to store, creating only 19 GB of metadata in the form of manifests. But because it points at the original NetCDF files, which must therefore be kept around, the full storage cost for the virtual method is 6.5 TB ($149), not 19 GB. Whether you count the NetCDF copy for the other methods depends on your policy. Many organizations choose to keep the originals regardless, for example, to support existing workloads that rely on file downloads. The TCO calculator below lets you count it either way with its NetCDF copy toggle.
Ingestion Cost and Latency
Ingestion was measured one cycle at a time on a single c7i.16xlarge, which is how a production system would work: a cycle lands, you ingest it.
The virtual ingest is 4x cheaper than any transform, because it only reads the HDF5 metadata of each file (about 26 MB of each 35 MB file, as it happens, since the chunk index is scattered through the file) and never decompresses or rewrites the actual data variables.
The three transforms are all in the same ballpark, bounded by reading and decompressing 232 GB of NetCDF and writing it back out. The Iceberg ingest was the slowest and by far the most memory-hungry: building the Arrow table for one lead time (22 million rows, 87 columns) and encoding it to Parquet peaked at 27 GB per process, which limited how many leads we could run concurrently on a 123 GB machine. We also hit operational friction that the tensor methods did not have. Ten concurrent writers appending to the same Iceberg table produced commit conflicts that needed retry tuning, and after 4,788 snapshots the per-commit metadata overhead was measurably slowing each append.
Latency matters too. Twenty-eight minutes is a long time when your forecast has a shelf life of six hours before the next cycle replaces it. For a team that needs the newest forecast available immediately, the virtual method’s six minutes (and the fact that it does not need to wait for the whole cycle before a file becomes queryable) is a real advantage.
There is undoubtedly room to optimize all of these ingestion pathways further, but the basic conclusions—virtual is cheaper than native, ingestion is I/O bound—are unlikely to change.
Query Cost and Performance
This is the most important chart because, at least in our cost model, it’s the dominant driver of divergent costs between the different approaches.
Times are median wall-clock seconds on the m7i.4xlarge instance, reading directly from S3 with no other activity on the host.
Digging into Q1, it also helps to understand how many bytes each method read to satisfy the point query: 2 MB for the timeseries layout, 83 MB for virtual, 315 MB for Iceberg, 6.1 GB for download, and 8.5 GB for FUSE, all to return 12 KB. In other words, Q1 suffers from strong read amplification, but the severity is strongly dependent on the storage method. The number of S3 requests behind each point query varies even more: 3 for the timeseries layout, about 170 for download, 300 for the dataloader layout, 2,500 for virtual, 5,500 for Iceberg, and 16,600 for FUSE. We measured these with CloudWatch, and they matter for the bill, as we’ll see below.
The native Icechunk layouts beat every other method on every query, and every Icechunk method beats the Iceberg table by at least an order of magnitude.
Let’s dig into why.
- Iceberg. The tabular layout gives DuckDB no way to find the bytes it needs without scanning: Q2 needs about 150 MB of values.
latandlonare not partition columns, so DuckDB reads the coordinate columns of all 1,368 files touching the two cycles at a few MB/s, and ends up over 10x slower than simply downloading those files. Hilbert sorting helped Q1 and Q2 by 4-5x but cannot help Q3, and a fresh DuckDB process pays another 50-75 s to plan over 19,152 files before its first query. - Download. Downloading 6 GB to read 12 KB is absurd, but a few hundred concurrent GETs saturate the instance’s 6 Gbps network, so Q1 “only” takes 12 seconds. The method is bandwidth-bound, and on the dataloader query, moving whole files at line rate is actually a little faster than issuing the virtual repository’s 100,000 small requests.
- FUSE. By far the worst method. Mountpoint pays about a second of overhead per file opened, because HDF5’s scattered metadata reads keep restarting the prefetcher, so Q1 (171 files) takes over four minutes however few bytes it needs. If you are wrapping legacy file-based code around S3, just download the files instead.
- Virtual Zarr / Icechunk. One second for a point query, with no data copied, is a great result: Icechunk’s Rust I/O layer issues the roughly 2,500 small S3 requests behind it very quickly. Those requests are not free, though, as the cost model below shows.
- Native Zarr / Icechunk. Each layout wins the queries it was designed for by a wide margin (three chunk reads for a timeseries Q1; training samples at 0.95 GB/s for the dataloader, 7x faster than virtual and 30x faster than Iceberg) and is the slowest Icechunk variant on the query it was not designed for. Even so, both beat virtual on all three queries, because their chunks need far fewer requests than 40 KB HDF5 chunks.
Total Cost of Ownership (TCO) Calculator
To turn these measurements into a bill, we need an assumed workload. Our baseline is a mid-sized team serving forecasts and training models:
| Workload | Frequency |
|---|---|
| Q1 point series | 10,000 per day |
| Q2 regional stats | 4 per day (once per cycle) |
| Q3 training epoch (16 samples per job) | 1 per day |
| Ingest new cycles | 4 per day |
Query cost has two parts.
Compute is priced as instance-seconds at the m7i.4xlarge rate, serialized, with no idle time.
This is generous to the slow methods; a real service pays for idle capacity.
S3 requests are priced at list rates ($0.0004 per 1,000 GETs, $0.005 per 1,000 LISTs) using the request counts we measured for each query with CloudWatch.
Storage is the method’s own bytes plus (in the “with NetCDF” column) the 6.5 TB of originals.
The interactive widget below allows you to tune the workload for your own anticipated usage pattern. A team that serves millions of point queries per day would want to increase Q1. An ML shop would want to increase Q3.
Conclusion
Iceberg costs 5x more than virtual Icechunk and nearly 8x more than the native timeseries layout at the baseline workload. Nearly all of the difference is the point query: 28 seconds of compute, 10,000 times a day, is 77 instance-hours a day, and the 5,500 range reads behind each one add another $670 a month in request charges. The storage and ingest bills are comparable across the transformed methods.
S3 request charges change the ranking. On compute alone, virtual Icechunk would be the cheapest method at this workload. But a virtual point query issues about 2,500 GETs against 40 KB HDF5 chunks, and 300,000 point queries a month turns that into $300, more than the virtual method’s storage and ingest combined.
The native timeseries layout is the cheapest option at the baseline workload, with or without the NetCDF copy. Three requests and a 100 ms per point query means its query bill is essentially zero, so its cost is almost entirely the fixed storage and ingest. Virtual only wins below about half our baseline query rate, where its cheaper ingest and negligible storage outweigh the request bill.
As the query workload grows, the gap widens. At 10x the baseline, the timeseries layout costs $461 a month against $3,871 for virtual and $25,853 for Iceberg. A team that wanted both native layouts (fast point queries and fast training) would pay two storage and two ingest bills, about $600 a month at 1x with the copy. This tradeoff should be quite appealing, and, in fact, we see many teams that choose to store the same data twice with different chunking schemes. (It’s what we do with our own Icechunk ERA5 dataset.)
Downloading whole files is a real option at low query rates, and it is 3x cheaper than Iceberg here. It just scales terribly: every point query moves 6 GB.
FUSE should not be used for this, at any workload.
This study settled the question we set out to answer: what is the most cost-effective system for storing and querying weather forecast data? We’ve shown that tensor-native storage on Icechunk beats a well-built Iceberg table on every query we tried, by one to two orders of magnitude, at a lower or comparable storage and ingestion cost. It also gave a more nuanced answer than we expected on which tensor layout to use: virtual for cheap, native for fast. As always, chunk shape is a key determinant of which queries are most performant.
What’s Next?
These results are an important and useful baseline for what you can do with a very simple architecture: just S3 and EC2, with no other persistent data, caches, or stateful services running. Any more complex system would sit on top of this baseline.
But there is clearly much more work to do.
-
More thorough comparison of database engines. We used DuckDB for the table side because it is fast, free, and reproducible. A warehouse like Snowflake, or a distributed engine like Spark, would change the request-concurrency picture for Iceberg, at a price. On the tensor side, Zax-SQL, our new SQL engine for arrays, would let us run the same SQL against both formats.
-
Ways to improve Zarr and Icechunk. While these results are already favorable towards Zarr and Icechunk, the experiments revealed several places where our stack could do better:
- Icechunk always issues a distinct S3 request for each virtual chunk, even when they’re in the same file. Intelligent coalescing of requests could lead to significant performance gains and reduction in I/O costs, similarly to how Zarr sharding works today.
- Even with the dataloader layout, we aren’t reaching either network or CPU saturation, indicating better performance is possible.
- We also found some sub-optimal behavior in sharded reads and a selection pattern that falls back to a slow indexer in zarr-python.
We’re actively working on resolving these issues.
-
Time to first forecast. We ingested whole cycles, and didn’t focus on making the forecasts available as fast as they were generated. This is probably the wrong assumption for latency-sensitive applications like energy trading. A more thorough comparison would take this dimension into account.
-
Operational overhead. Building a production-grade operational system is not as simple as firing up some EC2 nodes. Teams have to consider deployment, maintenance, security, stability, resilience, auto-scaling, and myriad other functional and non-functional requirements. That’s where the build-vs-buy calculus starts to change, which is why Earthmover offers fully managed compute services to help scientific data teams focus on their product rather than infrastructure. Because Earthmover Compute is built on the foundation of Icechunk and Zax, customers know they’re getting a best-in-class architecture which delivers the best possible value and performance for weather data workloads.
We’ll continue to research this important topic and will be publishing our findings as we go.
Reproducibility
Everything here is open-source and reproducible.
The ingestion code, queries, benchmarks, and cost model are in the open-source wxtco package, and the raw measurements are in the repository alongside the commands that produced them.
We’d love your feedback on anything and everything. Do you see any ways we could optimize any of the storage layouts? Know some tricks to make the DuckDB queries faster? We want to hear it.
CEO & Co-founder