Transfer-Efficient Data Processing in Disaggregated Systems
Abstract
Disaggregated analytics systems separate compute, memory, and storage to improve elasticity and resource utilisation, making network transfer a central bottleneck. Existing systems typically access remote data at fixed coarse granularities, transferring entire files, columns, or column-chunks even when queries process only small subsets of values. This mismatch between query selectivity and transfer granularity creates substantial unnecessary transfer. To address this problem, we propose making remote access granularity a first-class optimisation concern in disaggregated analytical systems. We propose a staged query execution approach that interleaves query evaluation with fine-grained remote loading using random-access storage layouts and intermediate selection vectors to fetch only the necessary data for further processing. By leveraging ideas of composable systems, the approach preserves compatibility with existing vectorised in-memory execution kernels. Using a prototype integrated with the BOSS composable DBMS, we show that fine-grained staged loading significantly reduces transferred data and improves runtime for selective analytical queries. We also show that memoisation allows lazy-loading overhead to converge to zero as an increasing fraction of data is accessed on the compute node. These results suggest that future disaggregated analytical systems should dynamically adapt remote access granularity during query execution to achieve transfer-efficient analytics.