Blog

Stay updated with our latest news and announcements.

Insights

[Summary] Integrating Distributed SQL Query Engines with Object-Based Computational Storage

Jeoungahn ParkJeoungahn Park (in SC25)


In modern large-scale data processing analytics systems, a disaggregated architecture that separates compute nodes from storage nodes is common. This architecture offers high scalability and manageability, but during analytics it retrieves entire data from remote storage — even when only a small fraction of data is needed — causing excessive data movement that degrade performance.

 

Existing object stores(e.g., AWS S3, MinIO) support only limited operations(e.g., Filtering, Column selection) to be processed on the storage side, while complex operations(e.g., JOIN, Aggregation) must still be performed on compute nodes leaving network bottlenecks unresolved.

 

This paper proposes a Presto-OCS connector that connects the Presto query engine with Object-based Computational Storage(OCS) which integrates a SQL query engine internally. By pushing complex operations to the storage, it minimizes data transfer over the network and significantly improves query performance

 

Key Concepts and System Design

Presto-OCS is designed to fully exploit the advanced pushdown capabilities of OCS without compromising Presto’s modular architecture. It integrates seamlessly into Presto’s existing execution framework by extending its Connector Service Provider Interface (SPI), enabling complex operators to be offloaded to storage while preserving compatibility with standard query pipelines. As shown in Figure 1:

 

Figure 1: Execution flow of the Presto-OCS integrated architecture

 

  • Pushdown Operator Determination: During query optimization, the connector automatically identifies data-reducing operators — such as filters, aggregations, and sorts — that can be pushed down to OCS. It leverages statistics from the Hive Metastore (e.g., min/max values, row counts) to estimate selectivity and prioritize operators with high data reduction potential.
  • Standardized Plan Translation: Identified operators are translated into Substrait Intermediate Representation (IR) — a cross-engine query plan format — and delivered to OCS to ensure seamless interoperability and in-storage execution.
  • In-Storage Execution: OCS receives the Substrait plan via gRPC and executes the operators using its embedded SQL engine. It returns results in Apache Arrow format. This minimizes network traffic by transmitting only processed output, not raw data.
  • Post-Processing: Operators not eligible for pushdown — such as complex joins, window functions and user-defined functions — are executed by Presto workers as usual. The system ensures full SQL semantics are preserved, with workers processing pre-computed data from OCS.

 

Performance Evaluation

In real-world High-Performance Computing (HPC) and OLAP workloads, Presto-OCS overcomes the limitations of traditional object storage and delivers significant performance gains: 


Figure 2: Execution time comparison with progressive query pushdown.

 

  • Laghos: Presto-OCS with complete pushdown of all operators achieves 2.25× faster execution and 99.99% reduction in data movement compared to filter-only pushdown.
  • TPC-H Query 1: On aggregation-heavy workloads, Presto-OCS delivers 4.07× speedup and 99.7% reduction in data transfer compared to filter-only pushdown.
  • Synergy with Compression: When combined with compression algorithms (Snappy, GZip, Zstd), it outperforms filter-only pushdown on compressed data by 1.36–1.39×.

 

Aggregation and sorting consistently improve performance, but complex column projections may hurt it — highlighting the need for smart pushdown strategies.

 

Conclusion

Presto-OCS is an innovative approach that eliminates network bottlenecks in disaggregated data analytics system by offloading data-reducing operators to storage. It fully leverages OCS’s advanced storage-side computation capabilities while preserving Presto’s architecture. Presto-OCS especially effective in HPC and large-scale OLAP workloads where network bottleneck is the primary performance constraint.


The Publication : Link



Popular Insights

Previous Next List