Benchmarks

Sail at Scale: Derived TPC-DS at 1TB and 3TB

All 99 derived TPC-DS queries in 444 seconds at 1TB and 1,289 seconds at 3TB on eight nodes, with shuffle streamed over Arrow Flight and never written to disk.

8 min read Oct 2026

Over two decades ago, Hadoop pioneered open-source distributed processing by following the ideas in Google’s MapReduce paper. The design was disk-centric. A computation was modeled as a sequence of map and reduce stages, and the data exchanged between them went through a distributed file system such as HDFS. Spark challenged that model with Resilient Distributed Datasets (RDDs), which keep intermediate data in memory when possible, track lineage, and recover a lost partition by redoing the computation. The shuffle between stages stayed disk-first, though. Shuffle data is written either to an executor’s local disk or to a remote shuffle service that is itself disk-based.

Rethinking Distributed Compute

When we started Sail, we envisioned another shift in distributed compute, toward doing even more of the work in memory. In a cloud-native environment it is easy to provision abundant compute for a short period, and network bandwidth inside a data center can be higher than what disks deliver. On those observations we built the first version of Sail’s distributed architecture in November 2024, with the Arrow Flight protocol carrying shuffle data between workers. Earlier this year we redesigned the control plane of our clusters to make the driver and workers more modular. The results in this post are the latest evidence of its effectiveness.

Benchmark Results

Sail completes all 99 queries of the derived TPC-DS benchmark in 444.3 seconds at 1TB and 1,288.8 seconds at 3TB, on eight r8gn.8xlarge workers that never write shuffle data to disk.

Each query ran in sequence and was timed end to end from the Spark Connect client, so every number includes planning, execution across the cluster, and the transfer of the result rows. The chart and table below give the totals. The two charts after them show the time for each query, first at 1TB and then at 3TB. The full setup and configuration are in the appendix.

Total query time, all 99 queries
02505007501,0001,2501,5001TB444.3 s3TB1,288.8 sTotal query time (seconds)
Metric1TB3TB
Total Query Time444.32 seconds1,288.77 seconds
Median Query Time2.62 seconds7.80 seconds
Queries Under 10 Seconds90 of 9969 of 99
Longest Queryq23, 38.12 secondsq85, 137.58 seconds

At 1TB, 78 of the 99 queries finish in under five seconds and the median query takes 2.6 seconds. The two longest, q23 and q14, take 38.1 and 33.8 seconds.

Execution time per query, 1TB
010203040Query time (seconds)q1q5q10q15q20q25q30q35q40q45q50q55q60q65q70q75q80q85q90q9538.1 s

At 3TB, the median query takes 7.8 seconds and 69 of the 99 finish in under ten. Three times the data took 2.9x the time overall. The longest are q85, q23, and q14 at 137.6, 106.1, and 92.1 seconds.

Execution time per query, 3TB
0255075100125150Query time (seconds)q1q5q10q15q20q25q30q35q40q45q50q55q60q65q70q75q80q85q90q95137.6 s

Arrow Flight Shuffle at Scale

The scalability of Arrow Flight shuffle may be counterintuitive at first glance. For two stages that each have p partitions, there are p × p logical shuffle streams between them. A complex workload with dozens of stages and a few hundred partitions, which is typical of the TPC-DS queries, can easily have more than a million logical streams active in the cluster. That does not mean the cluster needs that many network connections. Arrow Flight is built on gRPC and HTTP/2, and thanks to gRPC multiplexing, one TCP connection per pair of workers is enough.

In our experience a single connection per worker pair does well even for terabyte-scale jobs. We also support configuring the connection count for the Arrow Flight client to further boost performance if a single connection has a bandwidth limit in your settings. For example, in AWS, one TCP connection for EC2 has an upper limit of 5 Gbps by default if advanced capabilities such as Elastic Network Adapter (ENA) are not configured.

Sail is not alone in this design. Arrow Flight shuffle has recently come up in the DataFusion Comet and Riffle communities as well. We are honored to have started proving this direction two years ago.

Memory First, Not Memory Only

Pipelined shuffle over Arrow Flight completes the horizontal scaling story for memory-centric analytics. As we reported before, Sail is excellent at vertical scaling and runs 10x faster than Spark on a single machine. Enterprise applications may need to work through terabytes of data or more. With Arrow Flight shuffle you can provision a number of powerful nodes on a high-bandwidth network and never store intermediate data on disk. A task failure does require retrying the entire job under pipelined shuffle, but the cost of a retry is not what it was in disk-centric systems, because the job itself finishes in a fraction of the time.

Sail is not a memory-only engine, though. Many operators, sort among them, spill to disk when they need to. Sail also provides blocking shuffle backed by remote storage and by Apache Celeborn. If you operate in a resource-limited environment, or node preemption is frequent, blocking shuffle keeps a failure more contained. That is how Sail serves as one compute engine from your laptop to the cloud, while leaving you the choice of how to trade performance against resiliency.

Maturing a Distributed Engine

A distributed compute engine does not mature overnight, and looking back at the journey we are thrilled with how far it has come. The numbers reported here include our recent progress on bulk task scheduling, Arrow Flight transport configuration, and the query optimizer. All of it is available to our open-source users and to LakeSail Platform customers, so that production workloads at scale run faster and more robustly.

Getting Started

Install the pysail package and point your existing Spark SQL and DataFrame code at a Sail server. The Getting Started guide covers a single machine, and the Kubernetes deployment guide covers a cluster like the one in this post. You can follow along on GitHub or by joining our Slack Community.

Managed Sail in Your Cloud

Want to run Sail with managed infrastructure? The LakeSail Platform offers fully managed Sail in your own cloud account, with built-in governance, observability, and enterprise controls. Get started with a free trial and see how it can improve performance and reduce costs for your team. You can also schedule a free call with one of our engineers to ask any questions you might have.

Appendix

A. Benchmark Setup

TPC-DS is an industry-standard decision-support benchmark of 99 SQL queries over a retail data warehouse, covering multi-way joins, window functions, rollups, and multi-step common table expressions. We run a derived version using data of the standard schema and the SQLBench-DS query set, timed directly but without the official TPC audit. These are not official TPC-DS results.

We ran the experiment on the LakeSail Platform, which uses Amazon EKS to manage Sail servers and clusters. The Sail server and each worker occupy one EC2 instance apiece, and we fixed the number of workers so that resource allocation was predictable.

The data is stored on S3. We generated it with the Rust-based tpchgen-rs generator, with our own modification to write partitioned data. TPC-DS generation is still work in progress upstream, so that partitioning logic is not generally available yet. The dataset is Parquet with Snappy compression and 128 MB row groups.

Scale FactorRaw DataParquet DataPartitionsFiles
1,0001 TB430.9 GB2562,065
3,0003 TB1.3 TB7686,161

Partitioning applies to the fact tables and the large dimension tables (catalog_sales, catalog_returns, store_sales, store_returns, web_sales, web_returns, inventory, and customer). The remaining small dimension tables have a single partition each.

  • Queries: the 99 queries of SQLBench-DS, derived from TPC-DS under the TPC fair-use policy; q14, q23, q24, and q39 are two statements each, so a full run is 103 statements
  • Hardware: eight worker nodes on AWS EC2 r8gn.8xlarge (32 vCPU Graviton4, 256 GiB memory, 100 gigabits per second network), 256 vCPU and 2 TiB in total; driver on r8g.xlarge (4 vCPU, 32 GiB)
  • Disk: one 16 GiB gp3 root volume per node (3,000 IOPS, 125 MiB/s); no instance storage and no separate volume for temporary files
  • Network: all nodes in one availability zone with ENA Express enabled; no placement group
  • Client: PySpark over Spark Connect from a notebook on an r8g.xlarge in the same availability zone
  • Timing: each statement timed at the client from submission to the arrival of its result rows, fetched with toArrow(); the best complete run with that client setup is reported at each scale
  • Metrics: client-side wall-clock time only; no memory, disk, or network telemetry was collected in this run

Sail ran with the settings below changed from their defaults. Most of them are marked experimental in Sail’s configuration reference. Three idle and timeout settings were also raised so that workers and the session stayed up between runs (cluster.worker_max_idle_time_secs and spark.session_timeout_secs to 7,200 seconds, cluster.task_stream_creation_timeout_secs to 300).

Show the settings
SettingDefaultThis run
optimizer.enable_join_reorderfalsetrue
optimizer.enable_window_topnfalsetrue
execution.batch_size8,192262,144
execution.default_parallelism0 (number of cores)256
cluster.worker_task_slots896
cluster.shuffle_backend.flight.compressionzstdlz4
cluster.shuffle_backend.flight.connection_count14
cluster.shuffle_backend.flight.initial_connection_window_size0 (transport default)64 MiB
cluster.shuffle_backend.flight.initial_stream_window_size0 (transport default)8 MiB
execution.file_listing_cache.typenoneglobal
execution.file_listing_cache.max_entries10,0000 (unlimited)
parquet.file_statistics_cache.max_entries10,0000 (unlimited)
Cache TTLs (execution.file_listing_cache.ttl, parquet.file_metadata_cache.ttl, parquet.file_statistics_cache.ttl)1,800 s7,200 s

Here are four methodology notes.

  • The cost-based join reorder and the window top-N optimization are both marked experimental. We plan to enable join reorder by default once the optimizer implementation stabilizes.
  • Parquet metadata and directory listings were cached between queries, per the cache settings above. Table data was read from S3 for every query.
  • q85 took 1.9 seconds at 1TB and 137.6 seconds at 3TB. The plan chosen at 3TB under the experimental join reorder is the likely cause, and we have not yet isolated it.
  • This run measured Sail only. No other engine was run on the same cluster.

B. Spark and Spark Accelerators

We were not able to find authoritative TPC-DS results for Spark and Spark accelerators at these scales. The reports available online contradict one another, so we cannot draw a definitive conclusion from them for a comparison. We do note that on similar hardware configurations, Sail outperforms JVM Spark and Spark accelerators on the TPC-DS 1 TB and 3 TB datasets in similar settings.

A few heavy queries weigh on the total time and narrow Sail’s overall performance gain. Spark and Spark accelerators benefit from a query optimizer matured over many years, and that could inspire future work in Sail for even better performance. We recommend running controlled experiments on your own large-scale production workloads, and we would love to hear about your findings.

Questions about your Spark setup?

Book 30 minutes with a LakeSail engineer. Real answers, no pitch.