Blog

Stay updated with our latest news and announcements.

Insights

[Summary] Improving SQL Join Algorithms for Distributed Systems: A Case Study of Compute Express Link-Based Multihost Shared Memory

Jaeyung JunJaeyung Jun (in IEEE Micro)


For over ten years, the idea of Distributed Shared Memory (DSM)—which provides multiple processes with a unified virtual address space—has been actively studied in distributed systems. However, maintaining consistency and coherence in such systems is costly, as it requires moving data across locally attached memories in distributed servers. This has led most software frameworks to adopt the shared-nothing memory model, where processes share data only through explicit network communication.


In recent years, a new interconnect protocol called Compute Express Link (CXL) has emerged, enabling physical memory sharing across multiple servers. Initially designed to expand per-server memory capacity or bandwidth, CXL has evolved to support memory resource pooling across servers. Although many studies have examined memory pool usage in data centers, none have effectively exploited the data-sharing potential offered by such pooling.


In this paper, we present a case study investigating how data-sharing capabilities can be leveraged in distributed systems. We focus specifically on Spark SQL, a framework for running data analytics tasks expressed in standard SQL across distributed environments. Our attention centers on optimizing equi-join algorithms under a shared memory model—an operation that typically incurs heavy data movement costs in traditional distributed systems. An equi-join merges two tables by matching rows based on specified key columns. Traditionally, this requires a repartitioning phase to split tables into smaller, disjoint subsets that can be processed in parallel. However, this repartitioning involves a data shuffle—an all-to-all communication step—that generates significant network and file system overhead. Despite these costs, conventional equi-join algorithms remain the default choice under the shared-nothing model.


The rise of shared memory in distributed systems, enabled by the new CXL interconnect, prompts us to reevaluate the efficiency of existing equi-join algorithms under this new paradigm. These algorithms use repartitioning to narrow the search space from the full table to smaller, local subsets, which each executor processes independently. But what if all executors could directly access the entire table through shared memory? In such a scenario, reducing the search space would no longer be necessary, making the repartitioning phase obsolete.

 


Motivated by this insight, we propose a novel equi-join algorithm specifically designed for distributed systems with shared memory. Our algorithm minimizes data movement while preserving data parallelism for concurrent execution—with no lock contention. The key contributions of this paper are as follows:


  • - Improving SQL join performance by designing a new join algorithm optimized for the shared memory model.

  • - Presenting a practical approach to adapting existing system software for shared memory deployment with minimal engineering effort.

  • - Building and testing a CXL-based multi-host shared memory prototype, and evaluating the effectiveness of our new join algorithm within a distributed system using this prototype.


The Publication : Link



Popular Insights

Previous Next List