Distributed & Scalable Oblivious Sorting and Shuffling
Nicholas Ngai, Ioannis Demertzis, Javad Ghareh Chamani, Dimitrios Papadopoulos
IEEE Symposium on Security and Privacy 2024 · Day 3 · Continental Ballroom 6
Overview
In an era where data privacy is paramount, traditional encryption alone often proves insufficient to protect sensitive information from sophisticated attacks. This talk, presented by Nicholas Ngai and his colleagues from UC Berkeley, UC Santa Cruz, and HK, delves into the critical challenge of ensuring data confidentiality not just at rest or in transit, but also during computation. The core problem addressed is side-channel leakage, where adversaries can infer sensitive data by observing patterns in memory access, control flow, or network requests, even when the data itself is encrypted or processed within secure environments like hardware enclaves.

Key moments
- 0:00 Introduction and motivation for oblivious algorithms
- 1:50 Oblivious algorithms in Signal and other applications
- 4:00 The scalability problem in existing oblivious sorting
- 4:20 Our distributed, scalable oblivious sorting and shuffling solution
- 4:50 Four key optimization layers for improved scalability
- 5:10 Iterating on Bonic sort and its fundamental limitations
- 6:40 Transitioning to D-Bucket sort for enhanced scalability
Distributed & Scalable Oblivious Sorting and Shuffling
Speakers: Nicholas Ngai; Ioannis Demertzis; Javad Ghareh Chamani; Dimitrios Papadopoulos
Conference: IEEE S&P
YouTube: https://www.youtube.com/watch?v=0rn07vJx-jI
Overview
In an era where data privacy is paramount, traditional encryption alone often proves insufficient to protect sensitive information from sophisticated attacks. This talk, presented by Nicholas Ngai and his colleagues from UC Berkeley, UC Santa Cruz, and HK, delves into the critical challenge of ensuring data confidentiality not just at rest or in transit, but also during computation. The core problem addressed is side-channel leakage, where adversaries can infer sensitive data by observing patterns in memory access, control flow, or network requests, even when the data itself is encrypted or processed within secure environments like hardware enclaves.
The presentation introduces novel advancements in oblivious algorithms, specifically focusing on sorting and shuffling, which are fundamental primitives for a wide array of privacy-preserving applications. The speakers highlight a significant gap in existing oblivious solutions: a lack of scalability to handle the massive datasets prevalent in modern distributed computing environments. Their work directly confronts this limitation, offering fully oblivious, distributed, and scalable sorting and shuffling primitives that achieve unprecedented performance and capacity.
This research is highly significant for the future of secure cloud computing and privacy-preserving data analytics. By enabling efficient oblivious operations on large, distributed datasets, the presented techniques pave the way for more robust and practical privacy guarantees in real-world applications ranging from private contact discovery in messaging apps to anonymous machine learning inferences and differential privacy mechanisms. The talk not only details the theoretical underpinnings but also provides a comprehensive evaluation demonstrating substantial performance gains and scalability across various hardware configurations.
Background
▶ Watch: Introduction and motivation for oblivious algorithms (0:00)
The necessity for oblivious algorithms stems from the inherent vulnerabilities of traditional encryption when data is actively being processed. While encryption secures data against passive eavesdropping, active computation can inadvertently expose information through side-channel attacks. These attacks broadly fall into two categories: query pattern leakage, where the frequency or distribution of client requests reveals insights into underlying data (e.g., a restaurant menu revealing popular items), and software side-channel attacks targeting hardware enclaves. Hardware enclaves, such as Intel SGX, are designed to create secure execution environments within untrusted hosts (like cloud servers), processing secret data in isolation. However, even within these enclaves, adversaries can observe control flow leakage (which if branch is taken) or memory access leakage (which array or map elements are accessed), thereby inferring secret data.
Oblivious algorithms are designed to mitigate these leakages by ensuring that their query patterns, control flow, and memory accesses are entirely independent of the secret data being processed. This deterministic behavior, regardless of input values, makes them a cornerstone for robust privacy.
Real-world applications increasingly rely on these techniques. Signal, a popular messaging platform, uses oblivious algorithms (based on Oblix and Snoopy) for private contact discovery, allowing users to find contacts on the app without revealing their full contact list to the server or vice versa. Other applications include:
- Anonymizing Google's Key Transparency to conceal client queries.
- Searchable encryption in databases like Proton Mail and MongoDB.
- Privacy-preserving Big Data analytics (e.g., OPAC).
- Differential privacy, used by Google, Apple, Microsoft, and Amazon for collecting and analyzing user data without identification, often requiring oblivious shuffle primitives (e.g., Google's Prochlo system and Privacy Sandbox initiative).
- Large Language Models (LLMs), where oblivious algorithms can conceal input queries to protect user privacy.
Just as non-oblivious applications build upon fundamental primitives, so do oblivious applications. Common oblivious primitives exist for search indexes, key-value lookups, and graph processing. Among the most fundamental building blocks are sorting and shuffling. While a significant body of work exists in this area, dating back to Batcher's 1968 paper and more recent works in CCS '22 and '23, a critical limitation has persisted: scalability. Most prior solutions have only been tested up to 4 gigabytes of data and rarely target truly distributed settings. This talk directly addresses this gap, providing the first fully oblivious, distributed, and scalable sorting and shuffling primitives capable of handling massive datasets.
Key Findings
▶ Watch: The scalability problem in existing oblivious sorting (4:00)
The core contribution of this work is the development of fully oblivious, distributed, and scalable sorting and shuffling primitives that overcome the limitations of prior art. The key findings demonstrate a significant leap in the practical applicability of oblivious computation:
- Unprecedented Scalability: The proposed solutions, particularly the dBucket sort, were evaluated with sort sizes up to 128 gigabytes in a distributed environment utilizing up to 64 Intel SGX enclaves. This represents the largest distributed evaluation of oblivious sorting to date, a substantial increase over the typical 4-gigabyte limits of previous works.
- Dramatic Performance Improvements:
- For oblivious sorting, dBucket sort achieved a nearly 7x throughput increase over the optimized Bonic sort implementation when run across 64 enclaves. Even with 128 GB datasets, a 5.4x speedup was maintained.
- For oblivious shuffling, the solution outperformed prior work like Aura Shuffle even in single-enclave settings, demonstrating a nearly 10x speedup with 64 enclaves and a 9x speedup on 128 GB datasets.
- Exceptional Distributed Scaling Behavior: Increasing the number of enclaves from 1 to 64 resulted in an almost 21x speedup for sorting and a 16x speedup for shuffling, showcasing highly efficient parallelization. This contrasts sharply with older algorithms like Bonic sort and Aura Shuffle, which barely break even with single-threaded performance at high parallelism.
- Efficient Multi-threading: Even within a highly distributed setting, the solution scaled well with multi-threading. Increasing threads per enclave up to 8 yielded a 6.3x speedup for sorting and a 6.6x speedup for shuffling.
- Practical Speed: The system can obliviously sort 2 gigabytes of data in less than a second and 128 gigabytes in less than a minute. Shuffling is even faster, processing 2 GB in just over half a second and 128 GB in a little over half a minute.
- Enabling New Applications: The scalable oblivious sort primitive was demonstrated to enable new higher-level primitives, specifically a scalable oblivious batched key-value lookup called Snoopy++. This application achieved 700,000 requests per second (QPS) for a 16 million object database, a substantial improvement over Snoopy's 4,000 QPS under similar conditions.
- Key Optimizations Identified: Microbenchmarks highlighted the individual impact of crucial optimizations: the Aura compact merge split yielded up to a 3.28x speedup, the XOR-based oblivious swap provided up to a 1.7x speedup, and the O(N) bucket routing achieved up to a 2.87x speedup in distributed settings.
These findings collectively address the long-standing scalability challenge in oblivious computation, making privacy-preserving technologies viable for large-scale, real-world deployments.
Technical Deep Dive
▶ Watch: Our distributed, scalable oblivious sorting and shuffling solution (4:20)
The journey to achieving scalable oblivious sorting and shuffling involved a multi-layered approach, encompassing algorithmic, network, and low-level assembly optimizations.
The initial work began with iterating on Batcher's Bitonic Sort, a widely used algorithm for oblivious sorting due to its deterministic behavior. Bitonic sort recursively sorts halves of an array and then performs a bitonic merge. The team applied multi-threading optimizations and batched swap requests between enclaves to mitigate network latency, resulting in what they believe is the fastest implementation of Bitonic sort to date. However, its fundamental limitations persisted: an inherent N log^2 N runtime complexity and poor spatial locality, leading to excessive network communication across enclaves in a distributed setting. This non-optimal scaling behavior necessitated a more radical approach.
To overcome these shortcomings, the focus shifted to Bucket Sort, leading to the development of dBucket Sort. The basic concept of oblivious bucket sort involves:
- Shuffling: Assigning a random number (key) to each element.
- Routing: Using a bucket butterfly routing network to route elements towards their final sorted positions based on these random keys. Each bucket holds elements according to bits of the random key (e.g., first two bits).
- Permutation: Randomly permuting elements within each individual bucket to achieve a fully random permutation.
- Final Sort: Applying any standard comparison-based sort (which is oblivious here because the input ordering is now based on random keys, not the original data).
The dBucket sort introduces three critical optimizations to enhance scalability and performance:
- Optimized Butterfly Network Arrow Routing: The original butterfly network design often results in writing pairs of buckets to memory locations different from where they were read, leading to significant memory thrashing. The optimization ensures that during routing, the two destination buckets are the same as the two source buckets. This minimizes cache line contention between threads, which is crucial for multi-threaded scalability, while preserving obliviousness due to its deterministic nature.
- Efficient Merge Split Operator: At the heart of the butterfly network is the merge split operator, which takes two buckets and outputs two buckets sorted by one of the bits of their random keys. Traditionally, a full oblivious sort (like Bitonic sort) would be used for this. However, the key observation here is that only a separation of zeros and ones (based on the bit) is required, not a full sort. The team leveraged a highly efficient oblivious compaction algorithm from Sassy et al. (CCS '22) for this step, significantly reducing computational overhead.
- O(N) Network Communication: In a distributed bucket butterfly network, the last
log Elayers (whereEis the number of enclaves) traditionally involve extensive cross-enclave network communication. The dBucket optimization recognizes that after the firstlog Elayers, the final enclave position for any given bucket is already determined. Instead of continuous network swaps, the algorithm performs the initiallog Elayers locally, then executes a single, all-at-once rearrangement step to send buckets to their correct enclaves. After this, the remaining layers of the butterfly network can be completed locally without any further network overhead. This reduces network communication to O(N), where N is the total number of elements, as each bucket is sent over the network at most once. This deterministic rearrangement preserves obliviousness.
Beyond algorithmic improvements, the team also addressed the network fabric and low-level assembly optimizations.
Encrypted MPI Layer: Any distributed system requires robust inter-node communication. For distributed oblivious computation with hardware enclaves, an additional layer of security is critical. The widely used Message Passing Interface (MPI), while fast and easy to use, lacks built-in security mechanisms. Existing solutions like TLS suffer from performance degradation due to serialization, especially in multi-threaded contexts, while DTLS presents usability challenges. Prior attempts to layer security on MPI were vulnerable to attacks like replay attacks. To address this, the researchers introduced a novel encrypted MPI layer designed to provide security, ease of use, and high performance simultaneously, satisfying all three properties for the first time in an Enclave-to-Enclave communication context. This layer was implemented using Pitch as the MPI layer and mbedTLS for cryptographic operations.
Low-Level Assembly Optimizations (Oblivious XOR Swap): The oblivious swap operator, which either swaps two values or performs a no-op in an oblivious manner, is a fundamental building block. Traditionally, this is achieved using x86-specific assembly instructions like CMOV or AVX2 equivalents like VPBLENDVB, which execute conditionally but obliviously at the assembly level. The paper proposes a novel usage of the XOR swap by making its second step conditional. This creates an oblivious swap that offers several advantages:
- Portability: It can be expressed in pure C, eliminating the need for architecture-specific inline assembly.
- Vector/Scalar Width Support: It supports all existing vector and scalar widths, including 8-bit and 512-bit operations (AVX-512), and even variable-length vectors on architectures like RISC-V, which
CMOV-based operations may not. - Compiler Optimization: The architecture-agnostic C implementation makes it easier for optimizing compilers to improve code efficiency, avoiding potential pitfalls of inline assembly.
The entire system was implemented in approximately 10,000 lines of C code on top of Open Enclave SDK for Intel SGX p2 hardware enclaves, running on Azure confidential computing machines with Intel Xeon CPUs.
Demo / Proof of Concept
▶ Watch: Iterating on Bonic sort and its fundamental limitations (5:10)
To illustrate the practical impact of their scalable oblivious sorting primitive, the researchers demonstrated its application in building a scalable oblivious batched key-value lookup system, which they named Snoopy++. This serves as a direct comparison and improvement over Snoopy, a prior work currently used in Signal's private contact discovery.
The original Snoopy, while effective for smaller datasets, showed significant performance degradation as the database size increased. Their testing indicated that Snoopy could achieve up to 130,000 queries per second (QPS) with a 2 million object database, but this throughput drastically dropped to just 4,000 QPS when handling 16 million objects.
Snoopy++ leverages the newly developed scalable oblivious sort with a straightforward construction:
- Combined Sorting: Client requests and the server's data are combined into a single set. This combined set is then obliviously sorted.
- Linear Scan and Copy: A linear scan is performed over the now sorted, combined set. During this scan, data corresponding to the requested keys is obliviously copied from the data records into the respective request records.
- Oblivious Compaction: Finally, an oblivious compaction operation is performed to extract only the completed requests, effectively removing the server's raw data and returning the results to the client.
The latency of this system is primarily determined by the scalability and size of the underlying oblivious sort. Therefore, a highly scalable oblivious sort directly translates into a highly scalable key-value lookup algorithm. The critical metric becomes whether the sort time can meet the target latency. For instance, if a target latency of 1,000 milliseconds (1 second) is desired for a database of 16 million elements, Snoopy++ can achieve this using just 32 enclaves.
The performance comparison is striking: while Snoopy was limited to 4,000 QPS for 16 million objects in their tests, Snoopy++, utilizing the same hardware, was able to achieve an impressive 700,000 QPS. This demonstrates a dramatic increase in throughput, making oblivious key-value lookups practical for large-scale, real-world applications. The paper details the intricacies of this construction, emphasizing that as long as sufficient hardware is available to meet the target latency, Snoopy++ offers an incredibly scalable and efficient solution for privacy-preserving batched key-value lookups.
Defensive Implications
▶ Watch: Transitioning to D-Bucket sort for enhanced scalability (6:40)
The advancements in distributed and scalable oblivious sorting and shuffling presented in this talk have profound implications for security defenders and organizations handling sensitive data.
- Embrace Oblivious Algorithms for End-to-End Privacy: Organizations should recognize that traditional encryption is insufficient for protecting data during computation. Adopting oblivious algorithms is crucial for mitigating side-channel attacks that can expose sensitive information even when processed within secure hardware enclaves. This work demonstrates that such algorithms are now practical for large-scale deployments.
- Leverage Hardware Enclaves with Scalable Primitives: The use of hardware enclaves like Intel SGX is a critical component of the presented solution. Defenders should continue to invest in and utilize these secure execution environments. However, they must pair them with scalable oblivious primitives to ensure that the security benefits extend to real-world data volumes without becoming a performance bottleneck.
- Prioritize Scalable Privacy-Preserving Solutions: The demonstrated ability to obliviously sort and shuffle up to 128 GB of data in less than a minute across 64 enclaves means that privacy-preserving analytics and operations are no longer limited to small datasets. Defenders should actively seek out and integrate solutions built on these scalable primitives for applications like private contact discovery, searchable encryption, differential privacy, and anonymous LLM inference.
- Adopt Secure Inter-Enclave Communication: The development of a novel encrypted MPI layer highlights the importance of secure communication between distributed enclaves. Defenders should ensure that any multi-enclave architecture implements robust, performant, and easy-to-use encrypted communication channels to prevent network-based side channels or data interception between trusted components.
- Optimize at All Layers: The success of dBucket sort stems from optimizations across algorithmic, network, and low-level assembly layers. This multi-faceted approach serves as a blueprint for defenders developing or evaluating secure systems:
- Algorithmic Improvements: Choose algorithms inherently suited for obliviousness and optimize their distributed behavior (e.g., O(N) network communication for dBucket).
- Network Fabric: Implement secure, high-performance communication protocols.
- Low-Level Code: Utilize portable and compiler-friendly oblivious primitives (e.g., XOR-based oblivious swap) to maximize efficiency within enclaves.
- Enable New Privacy-Preserving Applications: The Snoopy++ proof-of-concept for scalable oblivious batched key-value lookup opens doors for new privacy-preserving services. Defenders can now confidently design and deploy services that require efficient lookups over large, sensitive datasets without compromising user privacy. This could include private data matching, secure identity verification, and anonymous telemetry collection.
In essence, this research provides the tools and methodologies for building truly private and performant distributed systems, moving the needle from theoretical possibility to practical implementation for a wide range of defensive use cases.
Key Takeaways
- Addressing the Scalability Gap: Prior oblivious sorting and shuffling algorithms lacked the scalability for modern distributed computing, typically limited to ~4GB datasets. This work provides fully oblivious, distributed, and scalable primitives capable of handling up to 128GB across 64 enclaves.
- dBucket Sort with Key Optimizations: The core innovation is the dBucket sort, enhanced by three critical optimizations: a shuffled butterfly network for reduced memory thrashing, an optimized merge split operator using oblivious compaction, and an O(N) network communication strategy that sends data across enclaves only once.
- Novel Encrypted MPI Layer: A new, secure, easy-to-use, and performant encrypted Message Passing Interface (MPI) layer was developed for robust and private communication between distributed hardware enclaves.
- Portable and Efficient Oblivious XOR Swap: A novel XOR-based oblivious swap operator was introduced, offering enhanced portability (pure C implementation), broader support for vector/scalar widths, and better compiler optimization compared to traditional CMOV-based instructions.
- Significant Performance Gains: The solutions achieve substantial speedups: nearly 7x faster than optimized Bitonic sort for sorting and 10x faster than Aura Shuffle for shuffling in distributed settings (64 enclaves). They also demonstrate remarkable scaling, with up to a 21x speedup when increasing enclaves from 1 to 64.
- Enabling New Applications: The scalable oblivious sort enables new higher-level primitives, demonstrated by Snoopy++, a scalable oblivious batched key-value lookup that achieves 700,000 QPS for 16 million objects, a massive improvement over Snoopy's 4,000 QPS under similar conditions.
About the Speaker(s)
Nicholas Ngai, the primary presenter, is affiliated with UC Berkeley, where he conducted this research. He collaborated with colleagues Ioannis Demertzis, Javad Ghareh Chamani, and Dimitrios Papadopoulos from UC Santa Cruz and HK. Their collective work focuses on advancing the field of privacy-preserving computation, particularly in the context of distributed systems and hardware enclaves. This presentation highlights their expertise in designing and optimizing fundamental cryptographic primitives for real-world scalability and performance.
Reviews
Dr. Zero (Offensive Security Researcher) — MUST SEE
This work shatters the scalability barrier for oblivious sorting and shuffling, a critical bottleneck for real-world privacy-preserving computation. The dBucket sort, with its O(N) network communication and novel optimizations, delivers unprecedented performance and capacity for large datasets. This isn't just theory; it directly enables practical, high-throughput privacy for applications from contact discovery to LLMs.
Heather Calloway (CISO) — STRONG ACCEPT
This research presents a critical breakthrough in scalable oblivious computation, making privacy-preserving sorting and shuffling viable for large, distributed datasets. It effectively addresses a long-standing performance bottleneck, enabling organizations to implement robust end-to-end data confidentiality during computation, which has direct implications for regulatory compliance and business risk management. This changes how security leaders should approach privacy architecture.
→ Top-rated talks at IEEE Symposium on Security and Privacy 2024