8 Query engine object storage performance determines how fast analysts get answers. Modern distributed SQL query engines read data directly from object storage in open file formats such as Parquet and ORC, usually organized by an open table format. They have no storage of their own to hide behind, so every query depends on how quickly the object store can deliver data to dozens or hundreds of workers at once. A well-tuned query cluster on slow storage is a slow cluster. This article explains how distributed query engines read from object storage, which storage metrics matter, how to size storage and network for a query deployment and how file layout and caching affect performance. It applies to any engine that separates compute from storage and reads over the S3 API. For the broader architecture, see our hub on running a data lakehouse on on-prem object storage. How a distributed query engine reads from object storage When a query runs, a typical engine: Plans the query using table metadata from a catalog and the table format’s manifest files. Prunes partitions and files that cannot contain matching rows, using statistics stored in metadata. Splits the remaining files into units of work and assigns them to workers. Has each worker issue S3 GET requests, often ranged reads of specific column chunks, in parallel. Streams results back through the cluster to the coordinator. Two consequences follow. First, a single query can issue thousands of small ranged reads at once, so request handling matters as much as raw bandwidth. Second, metadata reads (catalog lookups, manifest lists, file listings) happen before any data moves, so slow metadata operations add latency to every query. The storage metrics that matter Aggregate read throughput For large scans, the limit is how many gigabytes per second the object store can serve to all workers combined. A useful rule of thumb is to estimate the throughput each worker can consume when CPU-bound on decompression and filtering, multiply by the number of workers, and make sure the storage cluster and network can deliver that number with headroom. Requests per second and small-read latency Columnar formats encourage ranged reads of a few megabytes or less. A query touching thousands of files may issue tens of thousands of GET requests. Storage that handles large sequential reads well but struggles with high request rates will throttle queries. Ask for performance figures at realistic object and range sizes, not only for large sequential transfers. Time to first byte Interactive dashboards run many short queries. For these, the time between issuing a request and receiving the first byte matters more than peak throughput. Low and consistent first-byte latency keeps short queries fast. Listing and metadata operations Older table layouts depend on listing directories to find files. Modern table formats reduce this by tracking files in manifests, but catalogs and engines still issue HEAD and LIST calls. Storage with slow listing at large object counts creates a hidden bottleneck. Consistency Engines expect that once a write completes, readers see it. Strong read-after-write consistency on the object store avoids failed queries and stale results, especially when ingestion and querying run at the same time. Sizing storage and network Start from the workload rather than from capacity: Concurrency: how many queries run at once at peak, and how many workers each uses. Scan volume: how much data a typical heavy query reads after pruning. Target response time: what analysts and dashboards expect. From those, estimate the read throughput needed. Then check three places where the pipe can narrow: Storage nodes: each node can deliver a certain throughput, limited by drives, CPU and network interfaces. Scale out until the aggregate exceeds the target. Network fabric: the links between compute and storage, including any spine or top-of-rack uplinks, must carry the full aggregate. Oversubscribed links are a common hidden cap. Load balancing: requests should spread evenly across storage endpoints. Capacity planning still matters, but for analytics clusters performance usually drives node count before capacity does. File layout: the biggest lever you control Storage speed cannot rescue a poor file layout. The most effective tuning usually happens in the table, not the hardware. File size: very small files multiply requests and metadata overhead. Many teams target files in the range of a few hundred megabytes for large tables, and run regular compaction to merge small files written by streaming ingestion. Partitioning: partition by columns that queries filter on, such as date, but avoid partitions so fine that each holds only a few tiny files. Sorting and clustering: sorting data within files by common filter columns makes min and max statistics more selective, so more files and row groups are skipped. Row group and page sizes: sensible defaults work for most workloads; very small row groups increase request counts. Compression: modern codecs trade CPU for less data read. On fast storage, lighter compression can be faster; on constrained networks, heavier compression can win. Open table formats help by tracking file-level statistics and supporting compaction and maintenance without rewriting whole tables. Caching Most query engines support some form of caching, and it changes the storage requirement: Metadata caching keeps catalog and manifest information in memory, cutting repeated metadata reads. Data caching on local worker SSDs keeps frequently read file segments close to compute. Hot dashboards may read mostly from cache, while ad hoc queries still go to object storage. Result caching at the BI or engine layer avoids re-running identical queries. Caching reduces load on the object store but does not remove the need for strong throughput. Cold scans, new data and cache misses after restarts all hit storage directly. On premises versus public cloud object storage In public cloud, query clusters read from a shared service whose per-prefix request limits and network paths you do not control, and data transfer between regions or out of the cloud adds cost. On premises, you control the network between compute and storage and can design for the throughput you need, but you are responsible for sizing it correctly. Many organizations running large, steady analytics workloads choose on-premises object storage for predictable cost and performance, and keep sensitive data in a known location. Testing before you commit Synthetic benchmarks rarely reflect real query patterns. A practical test plan: Load a representative copy of a real table, with realistic file sizes and partitioning. Run your actual heavy queries and dashboard queries at expected concurrency. Measure end-to-end query times, storage throughput, request rates and latency percentiles. Repeat with caching disabled to see raw storage behavior. Add storage nodes or workers and check that performance scales roughly linearly. Watch for errors and retries as well. Throttling responses or timeouts under load are early signs that storage or network is the bottleneck. Common bottlenecks and how to spot them Workers idle while waiting on I/O: storage throughput or request handling is the limit. High CPU on workers, storage underused: compute is the limit; add workers or reduce decompression work. Slow planning before any data moves: metadata or catalog operations are slow, or tables have too many small files. Throughput plateaus as nodes are added: look for an oversubscribed network link or uneven load balancing. Performance degrades over time: small files are accumulating; schedule compaction. Checklist: object storage for a distributed query engine Aggregate read throughput sized for peak concurrency with headroom. High request rates and low latency for small ranged reads. Fast listing and HEAD operations at large object counts. Strong read-after-write consistency. Non-oversubscribed network between compute and storage. Even load balancing across storage endpoints. Healthy file sizes, partitioning and regular compaction. Metadata and data caching configured on workers. Testing with real queries at real concurrency before production. Putting it together A distributed query engine is only as fast as the storage it reads from. The storage must serve high aggregate throughput, handle large numbers of small requests with low latency and answer metadata calls quickly, all over a network that does not become the bottleneck. File layout and caching do the rest. Get those right and on-premises object storage can deliver interactive analytics at petabyte scale. For the full architecture, see running a data lakehouse on on-prem object storage, and for moving existing data, see migrating from HDFS to object storage. Frequently asked questions What storage performance does a SQL query engine need? Enough aggregate read throughput for peak concurrency, high request rates for small ranged reads, low first-byte latency and fast metadata operations. The exact numbers depend on worker count and query patterns. Is throughput or latency more important for analytics queries? Both. Large scans depend on throughput; interactive dashboards depend on latency and request handling. Most deployments need both. Do small files hurt query performance on object storage? Yes. Small files multiply requests and metadata overhead. Regular compaction into larger files is one of the most effective optimizations. Does caching remove the need for fast object storage? No. Caching helps repeated queries, but cold scans, new data and cache misses still read from object storage. Can on-premises object storage match public cloud for analytics? Yes, when storage nodes and the network are sized for the workload. On premises also gives predictable cost and control over data location.