Est.

Small File Performance Problems in Object Storage

HTTP round-trips make tiny files expensive at scale.

Senior Writer · · 11 min read
Cover illustration for “Small File Performance Problems in Object Storage”
Distributed I/O · October 4, 2026 · 11 min read · 2,418 words

Object storage treats every file as its own HTTP round-trip, and that single design choice explains almost everything that goes wrong when small files pile up. There's a flat namespace underneath, PUT and GET and DELETE and LIST requests instead of filesystem calls, and no shortcuts for handling lots of tiny things quickly. A block device has IOPS. A POSIX filesystem has IOPS. Object storage doesn't publish an IOPS number because the concept doesn't map cleanly onto how it works: requests move through a distributed queue rather than a fixed-capacity disk channel, so the cost of touching an object is baked in rather than something you can tune away.

Think about what a "read" even means in this world. On a local disk, you seek, you buffer, you let the inode cache absorb repeated access. In object storage, a read is an HTTP GET. You can ask for a slice of a file using the Range header, which helps, but there's no seek, no buffer, no cache sitting underneath to make repeated small accesses cheap. Every single object is a transaction, full stop on infrastructure terms: request out, response back, overhead paid.

None of this is a bug some vendor forgot to patch. It's the toll charged for near-infinite horizontal scale and a low cost per gigabyte, and that toll is precisely why object storage is the right foundation for data lakes and long-term archives. The system is built to hold enormous amounts of data cheaply and scale sideways without limit. The trouble starts when workloads generating millions of small objects get layered on top without anyone accounting for the per-request tax that comes with each one. The architecture didn't fail. The workload just didn't read the terms and conditions.

How small files accumulate in the first place

Small files rarely accumulate because of a single mistake. They show up because modern data pipelines are built to produce them, one reasonable decision at a time, until there's suddenly ten million objects where ten thousand would have done the job.

Streaming is the biggest contributor. Spark Structured Streaming writes a new file on every micro-batch. Kafka consumers with short trigger intervals do the exact same thing, batch after batch, hour after hour. Parquet, the columnar format most of this data is written to, doesn't support extending a file in place. So instead of appending to something that already exists, every micro-batch has to create a brand-new file. Running that every few seconds for a year drives the file count into the millions.

Partitioning schemes add another layer. Partition by something high-cardinality, like a user ID or a transaction ID, and the result is millions of folders each holding one small file, like mailing a single brick in its own box. Even schemes that sound sensible on a whiteboard, like partitioning by date, hour, and region, can carve a dataset into thousands of partitions where each one holds almost nothing.

Default settings do their own quiet damage. Distributed engines often write one output file per task partition unless someone deliberately coalesces the output, and that default rarely gets touched because nobody notices the file count until it becomes a problem.

Trace the root cause back far enough and it's usually CDC pipelines, Kafka consumers, and continuous ETL jobs. All three generate small files by design, not by accident. The small file problem, in a lot of shops, is really just the bill arriving for "we adopted streaming." Nobody did anything wrong individually. The sum of a thousand sensible choices is a filesystem nobody wants to query.

The per-request overhead that makes small files expensive

Every object needs its own GET request, and that single fact turns a dataset of millions of small files into a dataset of millions of billable, individually-timed API calls. What should be a handful of fast sequential reads becomes a swarm of tiny transactions, each one carrying its own fixed cost, each one adding up.

S3 bills per request. Reading a dataset made of a few large files costs a handful of GETs. Reading the same total volume of data split across millions of small files costs millions of GETs, charging you, in effect, for the privilege of fragmentation.

The latency cost runs even deeper than the dollar cost. Before a query engine reads a single row of data, it has to open the file, read its Parquet footer, parse the schema, and pull out column statistics. Opening a file typically costs one round-trip just for the footer, plus another for every column chunk the query actually needs. Multiplied by millions of files, the engine spends more time knocking on doors than it does sitting down and reading anything behind them.

Throwing more compute at this doesn't fix it. Adding workers increases how many of these requests can run at once, but it doesn't shrink the per-object overhead itself. The tax is fixed per file. More workers just means paying the same tax faster, in parallel, like hiring more clerks to process the same mountain of tiny invoices instead of consolidating the invoices first.

Metadata inflation as a compounding bottleneck

File count creates a second, separate bill, and this one gets paid by the catalog rather than the wallet. Every file added to a storage system generates metadata that has to be tracked somewhere, and once file counts climb into the millions, that tracking becomes its own performance ceiling, independent of how much actual data sits inside those files.

Hadoop clusters show this in its starkest form. The NameNode keeps filesystem metadata in memory, and every file and every block costs a fixed amount of that memory no matter how small the file is. Loading enough tiny files in makes the NameNode run out of room, so Hadoop shops have historically had to reach for workarounds like federation just to keep the metadata layer breathing.

Object storage has its own version of this problem, just wearing a different outfit. Table formats like Iceberg track partitions through manifest files, and when a table is made up of millions of small Parquet files, those manifests grow right along with the file count. Reading the manifests to figure out what to query becomes its own bottleneck, before a single byte of actual data gets touched.

Storage efficiency comes down to IOPS: seek time, read time, data transmission time. Metadata operations eat into that same IOPS budget: every cycle spent resolving file locations is a cycle not spent reading data. A query against a table with millions of small files spins up an equivalent number of tasks, one per file split, and the task scheduler can choke on that volume before any real processing starts.

A useful way to think about where this is heading: small file count, metadata bloat, and catalog fragmentation function as three separate but simultaneous forms of table decay. Fixing file count alone and ignoring the other two just means the table keeps degrading from a different direction. Increasingly, production maintenance treats all three together, because picking off one symptom while leaving the other two to compound is a losing trade.

How AI training workloads worsen the small file problem

AI training data manages to hit every failure mode described so far, all at once, at a scale that makes the structural tax impossible to ignore. Computer vision datasets in particular are often made up of millions of individual JPEG or PNG files, which is close to the worst possible shape for object storage to handle efficiently.

Older HPC workloads were built around large files, and the parallel file systems supporting them were tuned for exactly that: push a lot of bytes through a small number of big objects. AI workloads invert that profile. They're defined by a huge number of relatively small files, a pattern those older systems were never designed to serve well.

Picture a training job pulling from millions of individual JPEGs. On paper, the storage system has more aggregate throughput than the job could ever need. In practice, the system spends most of its time just locating and opening files rather than moving data. Metadata lookups become the bottleneck, not bandwidth, which is a bit like owning a sports car and discovering the real delay is finding your keys.

Sequential read tuning doesn't touch this problem, because the bottleneck isn't the read, it's the lookup. And adding GPUs doesn't help either, because GPUs sit downstream of a bottleneck that lives entirely in storage. The constraint is how fast data can move from storage into GPU memory, not how much compute is available to chew on it once it arrives. GPU utilization drops directly when data can't be delivered fast enough, and idle GPUs sitting around waiting on files are one of the more expensive ways to watch money do nothing.

The scale at which this plays out at Meta makes the point hard to dismiss. The company stores exabytes of training data and trains thousands of models against datasets ranging from terabytes to petabytes, and even with hyperscaler-level resources, Meta didn't have enough storage capacity to keep all that training data local to the hardware doing the training. That constraint pushed Meta toward a disaggregated storage architecture, and from there toward building the Data PreProcessing Service, a dedicated tier that handles extraction, transformation, and batching into tensors specifically to stop data stalls from idling the GPUs. A company with Meta's resources building an entire service just to keep training data moving is a pretty clear signal that this problem doesn't get solved by tweaking a few settings.

Query fan-out: how small files multiply work at the execution layer

Per-request overhead and metadata inflation recur at the execution layer, where they combine into something closer to gridlock. A query engine facing a table of millions of small files doesn't just pay for each file's I/O. It spins up a separate task for every file, and the overhead of scheduling, starting, and tearing down each of those tasks can eat more time than the actual data processing does.

Frameworks like Apache Spark and Hadoop MapReduce are built with large files in mind. Hand them millions of small ones instead, and the result is a flood of map tasks, each carrying its own cost for scheduling, initialization, and context switching. Past a certain file count, the scheduler itself becomes the bottleneck, so throwing more compute at the job doesn't speed anything up. It just means more workers standing in line waiting for the scheduler to hand out assignments.

Every file access also carries its own open-and-close overhead at the filesystem level, adding latency that has nothing to do with how fast the underlying storage can move bytes. The storage could be lightning fast and the query would still crawl, because the delay lives in the sheer number of operations, not their individual speed.

Cloud accounts add one more ceiling on top. Azure Storage accounts, for instance, can hit throttling limits at the account or share level once IOPS or ingress and egress thresholds are crossed. A query against a small-file table can get rate-limited by the storage provider itself, even when the actual data volume involved is fairly modest.

Per-request cost, metadata inflation, and execution fan-out aren't three separate problems: they're the same structural tax showing up at three different layers, so a small file problem rarely degrades gracefully. Query performance doesn't slide downward. It falls off a cliff, because every layer from the scheduler down to the storage API is getting squeezed at the same time.

The three remedies and what they fix

Three remedies have become standard for small files, each targeting a different layer of the problem described above, and none of them is a universal fix on its own.

Compaction rewrites a pile of small files into fewer, larger ones after the fact. It directly addresses query fan-out and per-request cost, but it's a retroactive cure, cleaning up what's already accumulated without stopping new small files from being written tomorrow. Modern table formats have largely absorbed this as a built-in feature. Apache Iceberg handles it through rewrite operations. Another vendor manages file sizing at write time through bin-packing, applied across both of its table types. Iceberg's snapshot isolation means compaction can run in the background without blocking anyone trying to read the table at the same time. Cloud platforms have pushed this further into the realm of "just turn it on." Dremio's Open Catalog automates compaction, expiry, and cleanup once table maintenance is switched on. AWS S3 Tables ship with managed compaction built in. Snowflake automates compaction for Iceberg tables as well. Apache Amoro offers an open-source version of the same idea, running alongside the catalog and triggering compaction once configured thresholds are hit. Hand-building a custom compaction pipeline is quickly becoming a thing of the past for anyone on a cloud-managed setup. Compaction cleans up the mess, but the root causes, like writer parallelism and streaming batch frequency, still need separate attention.

Container and archive formats take a different approach: pack many small files into one larger object so the storage layer sees far fewer individual items. This tackles metadata inflation and per-request cost at the same time. Sequence files, Avro, and Hadoop Archives all let small files travel as records inside one larger container, with the format itself acting as the aggregation layer. IBM's patented approach (US Patent 11,436,189) formalizes the same idea: combine selected files into a single aggregation file represented by one inode, which cuts metadata operation overhead substantially for both archiving and retrieval. The tradeoff is that reading an individual record back out requires the reader to understand the container format, which is a manageable dependency with well-supported formats but a dependency all the same.

Ingestion-time batching is the only one of the three that stops small files from being written. It works by buffering data before it ever gets committed to storage, so the files landing in the object store are already the right size. Lengthening micro-batch intervals, coalescing Spark output partitions before a write happens, and sizing Kafka consumer commit intervals to match a target file size are all forms of this same idea.

Each remedy covers different ground: compaction cleans up after the fact, container formats shrink the file count at the storage layer, and ingestion-time batching stops the problem at its source. None of them replaces the other two, and all three point to the same underlying lesson: the fix has to match the layer where the tax is actually being paid.

Diagram: Three Layers, Three Taxes, Three Remedies. Visualizes: Visualize how the small file problem manifests as the same structural tax at three distinct layers, and how each of the three standard remedies targets a different one.

Sources

  1. Performance- and cost-efficient archiving of small objects
  2. Performance of Distributed File Systems on Cloud Computing Environment: An Evaluation for Small-File Problem
Filed underDistributed I/O

More in Distributed I/O