Institute of Parallel and Distributed Systems University of Stuttgart Universitätsstraße 38 D–70569 Stuttgart Bachelorarbeit Design and Implementation of a NUMA-Aware Cooperative Scheduler Jonas Krieg Course of Study: B. Sc. Informatik Examiner: Prof. Dr. Christian Becker Supervisor: Simon König, M.Sc. Commenced: January 23, 2025 Completed: July 23, 2025 Abstract Efficient task scheduling is essential for maximizing computational processing unit (CPU) utilization in parallel applications. A widely adopted strategy is work-stealing, where idle threads dynamically steal tasks from busy ones. However, modern multi-core systems are increasingly based on Non- Uniform Memory Access (NUMA) architectures, in which memory access latency varies depending on the physical proximity of memory to processor cores. In such systems, traditional work-stealing algorithms—which primarily optimize for memory locality or load distribution—can lead to performance degradation due to load imbalance or remote memory accesses. Despite this critical need for both NUMA-awareness and dynamic load balancing in modern systems, existing scheduling approaches rarely address these requirements simultaneously. Most current work-stealing schedulers prioritize either memory locality or load distribution, failing to capture the complex interactions between memory access patterns and workload balancing. This oversight limits their effectiveness in optimizing performance for real-world, NUMA-based workloads. The goal of this thesis is to design and implement a scheduling system that explicitly considers the effects of Non-Uniform Memory Access and aims to optimize task execution performance in multi-threaded environments. This is done by extending existing scheduling concepts by integrating multiple critical parameters—specifically NUMA-awareness and system load balance—into the scheduling and work-stealing process. In this work, we design and implement a novel hybrid scheduling approach, called NUMA-Load- Aware Hybrid Scheduler (NLHScheduler), which combines the strengths of two state-of-the-art work-stealing strategies. Specifically, the NLHScheduler integrates NUMA-locality awareness with dynamic system workload balancing in every scheduling decision. Initial experiments revealed that these two criteria can sometimes conflict, leading to suboptimal scheduling decisions. To address this, the NLHScheduler prioritizes NUMA-locality, as previous analyses have shown that memory locality has a greater impact on performance than load balancing alone. Additionally, we enhanced the initial task assignment mechanism to be both NUMA-aware and load-sensitive, further improving scheduling efficiency. To evaluate the effectiveness of the proposed schedulers, a custom benchmark framework was developed. With this benchmark we tested various workload scenarios, including balanced and imbalanced task distributions, as well as different scaling behaviors by varying queue lengths, the number of thief coroutines, and the system’s concurrency level. The evaluation compared median execution times, system throughput, and successful steal operations across different schedulers. All work-stealing strategies significantly outperformed a baseline round-robin scheduler, improving performance by an average of 38.88%. While the individual stealing strategies showed similar results, the NLHScheduler achieved the greatest gains, especially by improving execution time stability by 40.59%. This work highlights the potential of combining NUMA-awareness and dynamic workload balancing in task scheduling, and lays a foundation for further research into adaptive, performance-oriented scheduling techniques in NUMA architectures. 3 Kurzfassung Effizientes Scheduling ist entscheidend, um die CPU-Auslastung in parallelen Anwendungen zu maximieren. Eine weit verbreitete Strategie ist das Work-Stealing, bei dem inaktive Threads dynamisch Aufgaben von stark ausgelasteten Threads übernehmen. Moderne Mehrkernsysteme basieren jedoch zunehmend auf NUMA-Architekturen, bei denen die Speicherzugriffszeit von der physischen Nähe des Speichers zu den Prozessorkernen abhängt. In solchen Systemen kön- nen traditionelle Work-Stealing-Algorithmen — die hauptsächlich nach Speicherlokalität oder Lastverteilung optimieren — aufgrund von Lastungleichgewichten oder entfernten Speicherzugriffen zu Leistungseinbußen führen. Trotz der entscheidenden Bedeutung von NUMA-Awareness und dynamischer Lastverteilung in modernen Systemen adressieren bestehende Scheduling-Ansätze diese Anforderungen selten gleichzeitig. Die meisten aktuellen Work-Stealing-Scheduler priorisieren entweder die Speicher- lokalität oder die Lastverteilung und vernachlässigen dabei die komplexen Wechselwirkungen zwischen Speicherzugriffsmustern und Lastausgleich. Diese Lücke schränkt ihre Effektivität bei der Leistungsoptimierung von NUMA-basierten Workloads ein. Ziel dieser Arbeit ist es, einen Scheduler zu entwerfen und zu implementieren, der die Auswirkungen von nicht-uniformen Speicherzugiffen explizit berücksichtigt und die Performance von Aufgaben in multithreaded Umgebungen optimiert. Dies erfolgt durch die Erweiterung bestehender Scheduling- Konzepte um die Integration mehrerer kritischer Parameter — insbesondere NUMA-Awareness und System-Last-Balance — in den Scheduling- und Work-Stealing-Prozess. Im Rahmen dieser Arbeit entwerfen und implementieren wir einen neuartigen hybriden Scheduling- Ansatz namens NLHScheduler, der die Stärken zweier moderner Work-Stealing-Strategien kom- biniert. Konkret integriert der NLHScheduler die NUMA-Lokalisierungsinformation mit dynamis- cher System-Last-Balance in jede Scheduling Entscheidung. Erste Experimente zeigten, dass diese beiden Kriterien sich gegenseitig widersprechen können und so suboptimale Scheduliung Entschei- dungen hervorrufen können. Zur Lösung priorisiert der NLHScheduler die NUMA-Lokalisierung, da frühere Analysen gezeigt haben, dass Speicherlokalität einen stärkeren Einfluss auf die Leistung hat als die reine Lastverteilung. Zudem wurde der Mechanismus zur initialen Aufgabenverteilung um NUMA-Awareness und Lastsensitivität erweitert, was die Effizienz des Schedulings weiter verbessert. Zur Bewertung der Effektivität der vorgeschlagenen Scheduler wurde ein eigenes Benchmark- Framework entwickelt. Mit dieser Benchmark testeten wir verschiedene Workload-Szenarien, darunter balancierte und unbalancierte Aufgabenverteilungen sowie unterschiedliche Skalierungsver- halten durch Variation der Queue-Längen, der Anzahl der Dieb-Coroutinen und der System- Nebenläufigkeit. Die Evaluation verglich die medianen Ausführungszeiten, den Systemdurchsatz und die Anzahl erfolgreicher Steal-Operationen verschiedener Scheduler. Alle Work-Stealing-Strategien übertrafen einen Baseline-Round-Robin-Scheduler signifikant und verbesserten die Leistung im Durchschnitt um 38,88%. Während die einzelnen Stealing-Strategien vergleichbare Ergebnisse erzielten, erreichte der NLHScheduler die größten Verbesserungen, insbesondere durch eine Steigerung der Stabilität der Ausführungszeiten um 40,59%. Diese Arbeit zeigt das Potenzial der Kombination von NUMA-Awareness und dynamischem Lastausgleich in der Aufgabenplanung auf und legt die Grundlage für weiterführende Forschungen zu adaptiven, leistungsorientierten Scheduling-Techniken in NUMA-Architekturen. 5 Contents 1 Introduction 17 2 Theoretical Background 19 2.1 Non-Uniform-Memory-Access (NUMA) . . . . . . . . . . . . . . . . . . . . . 19 2.2 Cooperative Scheduling . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 20 2.3 Coroutines . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 21 2.4 Work Stealing . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 22 3 Related Work 25 3.1 Locality Aware Work Stealing . . . . . . . . . . . . . . . . . . . . . . . . . . . 25 3.2 System Load Aware Stealing . . . . . . . . . . . . . . . . . . . . . . . . . . . 27 3.3 Hierarchy Aware Stealing . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 29 4 Design and Implementation 31 4.1 The Runtime Environment . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 31 4.2 Baseline Scheduler . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 34 4.3 Stealing Strategies . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 35 4.4 Advanced Task Scheduling - NUMA-Load-Aware Hybrid Scheduler (NLHScheduler) 43 5 Evaluation 47 5.1 Benchmark . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 47 5.2 Metrics . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 49 5.3 Parameters . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 51 5.4 Test Environment . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 52 5.5 Results . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 53 6 Conclusion and Outlook 65 Bibliography 67 7 List of Figures 4.1 Sequence diagram of the interacting between the worker and the worker_pool, for the stealing process of “worker 1”. . . . . . . . . . . . . . . . . . . . . . . . . . 33 4.2 Visualized effects of work stealing on the task queues. Worker 3 stole task E from worker 2. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 34 5.1 Median execution time of the benchmark with increasing queue lengths. . . . . . 53 5.2 Median throughputs of the benchmark with increasing queue lengths. . . . . . . 54 5.3 Total steals and remote steal-ratio of the benchmark with increasing queue lengths. 56 5.4 Median execution times of the benchmark with increasing heavy tasks. . . . . . . 57 5.5 Median throughputs of the benchmark with increasing heavy tasks. . . . . . . . . 58 5.6 Total steals and remote steal-ratio of the benchmark with increasing heavy tasks. . 60 5.7 Median execution times of the benchmark with a balanced workload and increasing concurrency level. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 61 5.8 Median throughputs of the benchmark with increasing concurrency level. . . . . 62 5.9 Total steals and remote steal-ratio of the benchmark with an increasing concurrency level. . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 63 9 List of Tables 4.1 Example of the workers’ states and their distances from node 0, the thief’s node. . 41 11 List of Algorithms 4.1 General Stealing Algorithm . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 32 4.2 Round Robin Task Scheduling . . . . . . . . . . . . . . . . . . . . . . . . . . . 35 4.3 NUMA-Aware Stealing . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 36 4.4 System Load-Aware Stealing . . . . . . . . . . . . . . . . . . . . . . . . . . . . 39 4.5 Combining Stealing Strategies . . . . . . . . . . . . . . . . . . . . . . . . . . . 42 4.6 Advanced Task Scheduling . . . . . . . . . . . . . . . . . . . . . . . . . . . . . 44 13 List of Abbreviations API application programming interface. 36 CCD Core Complex Die. 20 CPU computational processing unit. 3, 19 DBMS Database Management System. 28 FIFO First-In-First-Out. 26 HPC high-performance computing. 66 I/O Input/ Output. 20 IOD I/O Die. 20 LIFO Last-In-First-Out. 26 NLHScheduler NUMA-Load-Aware Hybrid Scheduler. 3, 7 NLLASS NUMA-locality and Load-aware Stealing Strategy. 40 NUMA Non-Uniform Memory Access. 3, 17 OS operationg system. 20 PU processing unit. 29 RAM Random Access Memory. 52 UMA Uniform Memory Access. 19 15 1 Introduction Since the 1950s, with the advent of the first digital computers, the field of computer science has undergone an extraordinary transformation. Early systems were constrained by limited processing power, memory, and simplistic architectural designs. Over the decades, continuous advances in hardware and software have led to the development of highly parallel and efficient computing archi- tectures [HP11]. One of the most significant milestones in this evolution has been the widespread adoption of multi-core and multiprocessor systems [DB11; HM08], which enable the parallel execution of tasks and significantly increase system performance across a wide range of application domains, if they are utilized properly. Davis and Burns [DB11] conclude that advancements in multiprocessor hardware have outpaced the development of effective methods for utilizing such systems. They suggest that this imbalance may result in systems running more slowly, despite improvements in the underlying hardware. As parallelism has become a cornerstone of modern computing, new challenges have emerged regarding how to efficiently distribute and manage workloads across processor cores [HBV19]. A key aspect of this challenge is the design of robust and efficient task scheduling mechanisms [BL99]. These must ensure not only balanced workload distribution and high throughput but also take into account architectural characteristics such as memory hierarchy and access patterns [KMS95; RDK+00]. In recent years, Non-Uniform Memory Access (NUMA) architectures have become prevalent in high-performance systems [ZBM+24]. In NUMA systems, memory is partitioned into regions (or nodes), each local to a group of processor cores [Lam13]. While this design improves scalability and memory bandwidth, it introduces non-uniform memory access latencies: accessing memory local to a processor is significantly faster than accessing memory on a remote NUMA node [BZFK10]. Consequently, memory locality becomes a crucial factor in scheduling decisions, especially for data-intensive and parallel workloads [BS96; GMT23]. However, many existing task schedulers fail to address the complexities of NUMA-locality and instead rely on simple strategies such as round-robin or random assignment, which can lead to performance degradation due to increased memory latency and suboptimal processor utilization [BSB19; DWXL18; SM16]. A promising approach to tackle these challenges is cooperative scheduling [HH23], an increasingly popular model in modern asynchronous runtimes [Kum20; TPTS17]. Unlike preemptive scheduling, where tasks can be interrupted arbitrarily, cooperative scheduling relies on tasks voluntarily yielding control. By allowing tasks to yield control, cooperative runtimes can make more informed decisions about memory locality. This control can help reduce cross-node memory accesses. Additionally, this approach allows for fine-grained concurrency with low overhead and simplified context switch- ing [Bou06; CF20]. Tasks are commonly implemented as lightweight coroutines, which execute until reaching suspension points, enabling efficient multitasking. Work-stealing, a dynamic scheduling strategy where idle worker threads proactively “steal” tasks from others’ queues, is compatible with both cooperative and preemptive scheduling, though it naturally complements the cooperative model. It facilitates effective load balancing without the overhead and complexity of preemptive task interruption [BL99]. However, conventional 17 1 Introduction work-stealing implementations often overlook NUMA characteristics or focus exclusively on either locality or system load awareness [Kum20; OPW+12; TPTS17]. This thesis addresses these gaps by designing and implementing a NUMA-aware scheduler that integrates both aspects, locality and system load awareness. This is achieved by implementing both, static and dynamic scheduling techniques, tailored for coroutine-based workloads in cooperative environments. The static component initializes task distribution with respect to memory locality and system load, while the dynamic component extends work-stealing with NUMA-awareness to maintain load balance and optimize memory access during execution. By combining the principles of cooperative scheduling with advanced, architecture-aware work-stealing, this work aims to significantly enhance performance and scalability on modern multi-core, NUMA-enabled systems. The resulting scheduler that takes all these aspects into credit is called NLHScheduler. Chapter 2 lays the theoretical groundwork by examining NUMA architectures, cooperative schedul- ing, coroutines, and work-stealing. This background establishes the context and key concepts necessary to understand the challenges and design choices explored later. Chapter 3 surveys existing research on schedulers, highlighting how prior work has approached scheduling through NUMA- awareness, system workload-awareness or memory hierarchy-awareness, and identifies the limitations this thesis seeks to address. Chapter 4 provides a detailed account of the design and implementation of the proposed scheduler. It outlines the rationale for each design choice, describes the workings of the runtime environment, highlights the dynamic scheduling approach through three distinct work-stealing implementations, and presents the static scheduling method featuring an advanced task scheduling algorithm. Chapter 5 introduces a custom benchmarking framework built to evaluate the scheduler under realistic workloads. Using coroutine-based matrix multiplication tasks, the benchmark naturally creates conditions that reveal the behavior of the stealing mechanism and test the scheduler’s performance. Metrics such as execution time, throughput, steal count, and scalability under balanced and imbalanced workloads provide a comprehensive assessment of the system. Finally, Chapter 6 summarizes the thesis’ contributions, reflecting on lessons learned and exploring directions for future research into efficient, architecture-aware schedulers for modern parallel systems. 18 2 Theoretical Background This chapter introduces the core technical concepts that underpin the design and performance of modern, coroutine-based runtimes. It begins by describing the architecture of NUMA, which governs how memory is accessed across multiple processor cores and directly affects performance in parallel systems. Next, it covers the fundamentals of task scheduling, contrasting preemptive and cooperative models, with an emphasis on the latter as it is commonly employed in coroutine runtimes. The chapter then explores coroutines as a programming abstraction that enables asynchronous and non-blocking execution patterns. Finally, it examines the work stealing scheduling paradigm, discussing how it interacts with NUMA architectures and coroutine-based workloads. Together, these concepts provide the theoretical foundation for understanding and evaluating runtime systems in modern multi-core environments. 2.1 Non-Uniform-Memory-Access (NUMA) NUMA is an advanced computer memory architecture primarily used in multiprocessor systems, where memory is distributed across multiple processing units. In contrast to traditional Uniform Memory Access (UMA) systems, where all processors have uniform access times to the mem- ory, NUMA introduces a heterogeneous memory access model [Lam13; LPM+13]. In a NUMA system, each processor is associated with a portion of local memory. This combination of processor(s) and local memory is referred to as a NUMA node. The defining characteristic of NUMA is that memory access times vary depending on the proximity of the memory to the processor. Access to a processor’s local memory is relatively fast, whereas access to memory attached to a different NUMA node (remote memory) incurs higher latency. This non-uniform latency can significantly impact the performance of applications [CD24; LPM+13]. Therefore, minimizing remote memory access is essential—a practice known as NUMA-aware programming. To apply this effectively, the system’s specific memory hierarchy must be known [Lam13; LPM+13]. A NUMA node typically consists of one or more processor cores located on the same socket, along with a physically connected portion of the main memory. The overall system memory is thus partitioned across these nodes, resulting in a distributed memory architecture. The number of nodes in a system depends on its design and can range from just one, resulting in a UMA system, to several dozen [Lam13]. NUMA systems are commonly used in high-performance computing environments, data centers, and large-scale servers, where the need for efficient memory access and processing power is paramount [Lam13]. NUMA is a key architectural feature in many modern computational process- ing units (CPUs), such as AMD EPYC and Intel Xeon Scalable processors [AMD20; WR20]. For instance, a dual-socket system using Intel Xeon Scalable, e.g. Ice Lake, CPUs typically consists of two NUMA nodes—one per socket—each with its own memory and up to 40 CPU cores. Local memory access latency in such systems is usually around 80 nanoseconds, while remote memory 19 2 Theoretical Background access latency can be approximately 1.3× to 1.7× higher, i.e., 100–140 nanoseconds, depending on the system interconnect and load conditions [Lam13; LPM+13]. Another representative example is the AMD EPYC architecture, particularly the second- and third-generation models such as EPYC 7742 (Rome) and EPYC 7763 (Milan)[AMD20; WR20]. These processors consist of up to eight individual Core Complex Dies (CCDs), each housing 8 CPU cores. The CCDs are connected via AMD’s Infinity Fabric to a central I/O Die (IOD), which handles memory and Input/ Output (I/O) interfaces. While memory is physically centralized in the IOD, logical NUMA domains are established by associating specific CCDs with portions of the memory. As a result, the operating system maintains memory locality, which reduces remote memory access latency [AMD21; WR20]. A single-socket EPYC 7742 processor (64 cores) can be configured in multiple NUMA layouts via BIOS settings [AMD20]: • NPS1: One NUMA node per socket • NPS2: Two NUMA nodes per socket • NPS4: Four NUMA nodes per socket These configurations affect how the processor’s cores and memory channels are grouped and how memory access patterns are interpreted by the operating system. While a single NUMA node (NPS1) simplifies scheduling, it may result in increased memory contention and remote access latency. In contrast, finer-grained NUMA configurations (e.g., NPS4) provide better memory locality at the cost of increased complexity in NUMA-aware application design [CD24; WR20]. Accessing remote memory across CCDs typically incurs a latency penalty of up to 2× compared to local access [Lam13; WR20]. These examples clearly demonstrate that NUMA is not a theoretical concern, but a practical performance factor that developers must consider when designing and deploying software on modern server-class hardware [CD24; LPM+13]. 2.2 Cooperative Scheduling Scheduling is the task of assigning units of work, called tasks, to workers for execution. In typical system models, a worker corresponds to a thread of execution, often realized as an operationg system (OS) thread that is mapped to a hardware core or logical processor. A task represents a schedulable unit of work, such as a function, job, or process, and can vary in granularity depending on the abstraction level. Two fundamental types of scheduling strategies are commonly distinguished: preemptive and cooperative scheduling. In preemptive scheduling, the operating system or runtime system has the ability to interrupt a currently running task at almost any time in order to schedule another task. This mechanism requires hardware support, such as timer interrupts, and is typically implemented at the OS kernel level. Preemptive scheduling allows for more responsive and fair task handling, especially in multi-user or real-time systems, but also incurs overhead due to frequent context switching [TB14]. In contrast, cooperative scheduling relies on the voluntary yielding of control by tasks themselves. A task continues running until it explicitly yields execution, typically by calling a yield function provided by the scheduler. This model simplifies scheduling logic and reduces overhead, but comes with the risk of one task monopolizing the CPU if it fails to yield properly, thereby starving other tasks. Cooperative scheduling is commonly found in user-space runtimes such as event loops, e.g., 20 2.3 Coroutines Node.js, or coroutine-based systems [HH23]. In both models, the scheduler is still responsible for assigning tasks to workers, but the timing and mechanism of context switches differ significantly. Preemptive scheduling depends on external intervention, e.g., kernel timers, while cooperative scheduling delegates responsibility to the running tasks themselves. Nonetheless, many user-space runtime systems—such as Go’s scheduler and Scala’s ZIO—operate on top of preemptively scheduled OS threads. These systems multiplex lightweight, cooperative tasks (fibers or goroutines) onto OS threads, ensuring concurrency without starvation and enabling parallel execution across cores [Con22; LLC20]. 2.3 Coroutines One use case of cooperative scheduling is the controlling of coroutines. A coroutine is an abstraction of the subroutine, which is traditionally used in programming. A subroutine has a defined control flow: 1. Invocation 2. Execution 3. Completion A coroutine is a generalization of a subroutine that allows its execution to be suspended and resumed at specific points. Unlike a traditional subroutine, which has a single entry and exit point and executes linearly, a coroutine can pause its execution (suspend) and later continue from where it left off (resume), maintaining its internal state between invocations. The typical control flow of a coroutine can be described as follows: 1. Start: The coroutine is invoked and begins executing from its initial entry point. 2. Execution: • Suspension: At certain points (e.g., via yield or await), the coroutine can suspend itself, saving its local state. • Resumption: Control can later be passed back to the coroutine, resuming its execution from the last point of suspension. • Suspension • Resumption • . . . 3. Completion: The coroutine eventually finishes execution and returns a final result or exits normally. It is important to note that suspension and resumption are optional: a coroutine can behave like a traditional subroutine if it does not yield or suspend. However, its power lies in the ability to be paused and resumed multiple times, which makes coroutines particularly well-suited for asynchronous and event-driven programming [PLMA17]. Asynchronous programming is important because it decouples the initiation of a task from its completion. This decoupling allows programs to continue executing while waiting for long-running 21 2 Theoretical Background operations, e.g., I/O, network requests, or computations, which reduces idle time and improves overall system utilization. As a result, asynchronous programming has become a widely adopted paradigm in high-throughput and latency-sensitive applications such as web servers, distributed systems, and user interfaces [Myk24; SMMG22]. However, asynchronous programming introduces its own challenges. Traditionally, it relied heavily on callbacks—functions passed as arguments to be executed when an asynchronous operation completes. This approach often leads to deeply nested and hard-to-read code structures, a phenomenon commonly referred to as “callback hell” [EFGK03]. Coroutines offer a more structured and readable way to write asynchronous code. They allow a function to be suspended and later resumed, all while maintaining its internal state. This enables asynchronous code to be written in a linear, sequential, i.e. synchronous, style without blocking the thread on which it runs [PLMA17]. In a synchronous program, subroutines would be blocked during I/O or network latency. In contrast, with coroutines, the function can suspend execution at a yield or await point, allowing control to be returned to the runtime scheduler, which can resume the coroutine once the awaited operation completes. The advantages of cooperative scheduling become apparent when coroutines call other coroutines. When a coroutine awaits another, it suspends its own execution and voluntarily yields control, allowing other tasks to run in the meantime. This behavior aligns naturally with cooperative scheduling, in which tasks are only switched when they explicitly yield [HH23]. Thus, coroutine-based runtimes are typically designed around cooperative schedulers, which assign and resume tasks only when those tasks have willingly yielded control. This makes cooperative scheduling a natural and efficient fit for coroutine-based asynchronous programming. 2.4 Work Stealing Work stealing is a thread-based scheduling paradigm that enables an idle thread to take over, i.e. “steal”, work from another thread. Unlike traditional scheduling, where a central scheduler assigns tasks to workers, in work stealing the decision-making is distributed: each thread maintains its own task queue and performs local scheduling. If a thread becomes idle, i.e., its queue is empty, it becomes a thief and attempts to steal work from other threads’ queues. Thus, the scheduling responsibility lies primarily with the worker threads themselves. The work stealing process involves two decisions: 1.) Find a busy core from whom to steal, 2.) relocate the workload to the idle thread. By migrating work from busy to idle threads, the scheduler can balance the workload of the whole program, which leads to a more efficient memory usage and higher performance. [TPTS17]. Work stealing contrasts with work sharing, where tasks are proactively distributed by a centralized scheduler, or a spawning thread, to balance load ahead of time. Work sharing attempts to prevent load imbalance by spreading new tasks immediately, whereas work stealing reacts to imbalance after it occurs. While work sharing can achieve good balance under predictable workloads, it often leads to more synchronization overhead and contention for a central queue, especially on highly parallel systems [BJK+95]. In NUMA systems, work stealing is often the preferable approach due to its decentralized and local-first nature. Since each thread maintains its own queue, task stealing can be optimized to prioritize victims on the same NUMA node, thereby reducing the latency penalties of cross-node memory access [Kum20]. In contrast, work sharing may unintentionally assign tasks across NUMA boundaries, leading to frequent remote memory accesses and degraded performance. There are two main policies for work stealing: work-first and help-first [GBRS09]: In the work-first 22 2.4 Work Stealing approach, the executing thread runs the newly spawned task and allows the continuation to be stolen by others. In contrast, the help-first approach executes the continuation and leaves the newly spawned task to be stolen by other threads. The work-first policy favors scenarios, where stealing is rare. The thread executes it’s task itself and only allows the continuation to be stolen. This results in less synchronization and coordination with the other threads. The help-first policy on the other hand is designed to allow as many steals as possible, because spawned tasks are desired to be stolen directly. 23 3 Related Work Scheduling for multi-core systems—particularly for NUMA architectures—has been a longstanding topic of research. The central objective in this area is to improve application performance by maximizing data locality and achieving efficient load balancing. A major focus has been on the development of work-stealing schedulers, which dynamically redistribute tasks across processing units to address workload imbalances. Prior work in this field can be broadly categorized based on the type of information used to guide scheduling decisions. We identify three main strategies: • Locality-aware scheduling, which attempts to keep tasks close to their associated data. • System load-aware scheduling, which prioritizes balancing workload across cores. • Hierarchy-aware scheduling, which takes into account the static hardware topology such as cache levels and NUMA node structure. The following sections present key approaches and research contributions within each of these categories. Where appropriate, we discuss their assumptions, limitations, and relevance for modern coroutine-based runtimes. Cooperative scheduling is a scheduling paradigm that inherently favors the usage of coroutines, but is not limited to them. Hähnle and Henrio [HH23] show that a cooperative scheduler can be implemented in a way that it is provably fair without coroutines. Their work focuses on proving that a cooperative scheduler can be fair in general, so they do not focus on NUMA architectures, which is our main focus. 3.1 Locality Aware Work Stealing Yoo et al. [YHK+13] propose a locality-aware work-stealing scheduler that aims to improve cache efficiency by grouping tasks based on data locality. Their approach begins with an initial task scheduling phase, where tasks are assigned to workers in a way that attempts to exploit spatial locality. To do this, they construct a task dependency graph, where nodes represent tasks and edges indicate shared data or cache usage between tasks. Based on this graph, the scheduler tries to group tasks together that share common data and are therefore likely to benefit from being executed on cores that share caches. However, partitioning such a graph into optimal task groups, such that intra-group data locality is maximized and inter-group interference is minimized, is a task that requires heavy computing resources, unfit for a scheduler. Because of this, [YHK+13] resort to heuristics for partitioning, which introduces the risk of suboptimal groupings. Poor groupings can either lead to small task groups, causing high scheduling overhead, or overly large ones, which may overload individual cores and lead to load imbalance. In addition to grouping, the scheduler introduces a locality-preserving execution order within each group, based on two cache-related metrics: the intra-group sharing degree, which quantifies 25 3 Related Work how much a task shares cache lines with others in the same group, and the inter-group sharing degree, which captures sharing with tasks outside the group. Tasks with high intra-group and low inter-group sharing are prioritized to maximize cache reuse and minimize cache pollution. Task stealing operates at the granularity of whole task groups to maintain locality benefits. The stealing process follows the system’s cache hierarchy: a thread first attempts to steal from nearby cores (e.g., sharing L2 or L3 cache) before expanding outward. A key limitation of this work is its focus on single-processor systems with multiple cores and shared caches. As such, it does not account for NUMA-specific issues such as remote memory access costs across nodes. This reduces the applicability of their approach in distributed-memory systems, where memory locality is not solely determined by cache hierarchy but also by physical memory placement. Despite this, the paper remains relevant for two reasons: First, its use of locality-preserving metrics and cache-aware stealing demonstrates the potential performance gains of locality-aware scheduling strategies. Another approach focusing on work stealing in NUMA systems is proposed by Olivier et al. [OPW+12]. They present an extension of the OpenMP runtime that introduces NUMA-awareness into task scheduling. Their work is built upon the classic Cilk work-stealing model [BJK+95], which uses Last-In-First-Out (LIFO) dequeues to manage each thread’s local task queue. In this model, threads push and pop tasks from their own queues in LIFO order, while stealing is performed in First-In-First-Out (FIFO) order, i.e., thieves steal the oldest task from the victim’s queue. This approach helps preserve execution locality for spawned child tasks, which are often more tightly coupled, like we described in the previous work [YHK+13]. In the NUMA-aware extension, a thief thread selects a victim either randomly or based on the numerically closest thread ID. The intention behind choosing the “nearest” thread ID is to approximate memory locality. However, this heuristic proves problematic: thread IDs do not necessarily reflect physical proximity or memory domain boundaries. For instance, if threads are assigned to cores in a round-robin fashion across NUMA nodes, two sequential thread IDs may reside on different memory nodes, thereby violating the NUMA-affinity principle and leading to unnecessary remote memory accesses. This is not the case for any computers that we consider to test our implementation with, but still could be the case for other NUMA computers. Another shortcoming of their method is the lack of consideration for cache hierarchy. While their approach addresses NUMA to some extent, it treats the memory system as a flat hierarchy and does not exploit finer-grained locality within caches, e.g., L2 or L3 sharing between cores. This limits the scheduler’s ability to improve cache reuse or reduce interconnect traffic between cores and sockets. The authors acknowledge that, despite its simplicity and good baseline performance, work stealing in NUMA systems still faces limitations due to these blind spots in locality awareness. Their findings are supported by other works [ABB00; BFJ+96; CGK+07], which emphasize that while work-stealing scales well and reduces contention due to its decentralized nature, it must be extended with hierarchical and topology-aware stealing strategies to perform optimally on modern hardware. Blumofe and Leiserson’s foundational analysis [BL99] shows that work-stealing schedulers can significantly outperform work-sharing strategies, particularly due to reduced synchronization overhead and better load balancing. However, their original model does not account for memory or cache locality, which are crucial for performance on NUMA systems. Later research thus builds on their model by incorporating affinity-aware task placement and stealing strategies that consider NUMA domains or shared cache levels. Compared to these works, our approach aims to address these limitations by developing a scheduler that takes NUMA locality into account. Pufferfish [Kum20] is a work-stealing runtime that combines locality-aware and hierarchy-aware strategies to optimize performance in NUMA systems. Their work-stealing approach is designed to 26 3.2 System Load Aware Stealing minimize cache misses by ensuring that tasks are executed close to the data they operate on, which is critical in NUMA architectures. The key innovation in Pufferfish is the application of NUMA-aware work-stealing along with a hierarchical analysis of the system’s memory topology. This hierarchical analysis leverages the HCLib [KZC+14] library, which provides a generalized version of place, a concept introduced in languages like X10 [ESS05] and Chapel [CCZ07] for task locality. The place abstraction in HCLib helps determine the best location for a task within a NUMA cluster, improving memory locality. Additionally, HCLib enables asynchronous task execution through the async-finish paradigm. In this model, tasks are launched asynchronously, but the program waits for them to finish before continuing. However, the async-finish model has limitations. First, it requires the programmer to explicitly define the task placement, which introduces potential errors and reduces flexibility. Secondly, while async-finish allows for asynchronous task execution, the program flow remains fundamentally synchronous due to the need for synchronization points at each finish, which limits the ability to fully exploit concurrency. In contrast, our approach uses coroutines to guarantee a truly asynchronous control flow. Another important aspect of Pufferfish is its use of the help-first policy [BJK+95], which prioritizes helping other threads by stealing tasks from them. While this policy aligns well with the async-finish paradigm, it introduces unnecessary overhead for coroutines, as it can lead to more synchronization than needed, especially in a highly concurrent environment. In our scheduler, we aim to optimize the stealing policy to better suit coroutines by reducing synchronization overhead, which in turn reduces the overall complexity and improves performance. Furthermore, Pufferfish relies on the “NUMA memory manager” [Kum20] built using the libnuma library [Kle05] to manage NUMA-specific memory allocation. However, libnuma has a significant limitation: it provides only NUMA-specific functionality, which requires knowledge of the system’s topology beforehand. This static approach is problematic in dynamic environments, as it does not allow for flexible runtime adjustments. In contrast, our scheduler is designed to be more adaptable by performing runtime topology analysis, which allows the system to dynamically adjust to changes in the architecture. This approach provides more flexibility, enabling the scheduler to perform optimally on systems with varying NUMA configurations without requiring prior knowledge of the topology. In summary, while Pufferfish improves performance by combining locality-aware work-stealing with hierarchical system analysis, it still relies on static memory management and a synchronous control flow. Our approach builds upon the ideas in Pufferfish by using a coroutine-based runtime, reducing synchronization overhead, and enabling dynamic runtime topology analysis. This allows for more flexible and efficient scheduling in NUMA systems, particularly when handling unstructured parallelism with fine-grained concurrency. 3.2 System Load Aware Stealing Tzilis et al. [TPTS17] propose a system load-aware work-stealing strategy (SWAS), which aims to improve task scheduling by factoring in both the system load and task locality. The core idea of SWAS is to make work-stealing decisions based on the current load status of each core in the system. Each core has one of three possible states: Idle, Busy, or Don’t Know. The Don’t Know state is the default state for a core, indicating that no information about its load is available. To manage the system load information, each core stores a decentralized view of the load on all other cores. This view is represented as a set of bit vectors on each core, where each entry consists of two bits. These bit vectors track the load status of each core and are used to make decisions about work stealing. A core does not need to communicate directly with other cores to obtain this information. 27 3 Related Work Instead, it updates its own view of the system’s load based on the outcomes of work-stealing attempts. This design minimizes communication overhead, making the process of deciding where to steal work efficient, with only minor additional latency during the stealing attempt. In addition to load information, SWAS also incorporates locality awareness. Each core maintains a list of regions that define the organization of the system’s memory hierarchy. In a NUMA system, these regions could correspond to the different NUMA nodes, with the cores in each node categorized into distinct regions. The list of regions is sorted by proximity to the current core, meaning that cores that are closer in terms of memory access latency are considered to be in higher-priority regions for stealing work. When a core tries to steal work, it first checks the regions to identify the one that has the most work to be stolen, ensuring that the thief’s action is based on both load and locality information. However, while this strategy works well in a general multi-core system, it has some limitations when applied to NUMA systems. The main issue arises because, in SWAS, work-stealing is prioritized based on the system’s load distribution, not the locality of memory access. This means that when a core tries to steal work, it may choose to steal from a region with the highest load, even if it is far away in terms of memory access, e.g., on a different NUMA node. While this may balance the load across the system, it ignores the locality of the data, which is crucial in NUMA systems. The previous presented studies [Kum20; OPW+12] have shown that locality is the most important factor in NUMA systems because accessing memory on a distant NUMA node incurs higher latency. In contrast to SWAS, our approach incorporates a stronger emphasis on memory locality and NUMA-awareness in the task scheduling process. We aim to ensure that tasks are scheduled on cores that are closer to the data they operate on, minimizing memory access latency. Our scheduler does not rely solely on system load as the primary criterion for work-stealing. Instead, it uses a more dynamic approach, adjusting task assignments based on the actual topology of the NUMA system. By conducting runtime topology analysis, our scheduler works on any memory architecture, enabling better performance on systems with complex NUMA configurations. In summary, while SWAS focuses on load-awareness and locality through a decentralized bit-vector system, our approach goes beyond load balancing by combining NUMA-aware task scheduling with real-time system topology analysis, with load balancing. This ensures that locality is always a part of the stealing decision. Psaroudakis et al. [PSM+16] focus on the challenge of load balancing in NUMA architectures, particularly for memory-intensive tasks, and propose an adaptive scheduler designed to optimize task distribution for Database Management Systems (DBMSs). Their key insight is that work-stealing, which works well for load balancing in general, can actually decrease performance when dealing with memory-intensive workloads, where maintaining locality is more important than balancing the load across all cores. To address this issue, the Psaroudakis et al. propose an adaptive scheduler that dynamically switches work-stealing on and off depending on the workload’s characteristics. They identify an experimentally determined threshold that triggers the change in strategy. Specifically, when the workload is memory-intensive, work-stealing is disabled to avoid unnecessary overhead and memory access penalties, as stealing work from a distant NUMA node can introduce significant latency. On the other hand, for less memory-sensitive tasks, work-stealing is enabled to balance the load and improve performance. While this adaptive strategy may work well for database management systems, which typically involve a mix of memory-intensive and computationally-intensive tasks, it is not directly applicable to all types of parallel applications. The focus on DBMSs means that the solutions provided are tailored to that domain and may not address the broader performance requirements of other workloads, such as those involving unstructured parallelism or coroutines. The main limitation of Psaroudakis et al.’s work in the context of our research is that it focuses on disabling work-stealing for memory-intensive tasks, whereas our scheduler aims to optimize task 28 3.3 Hierarchy Aware Stealing placement and work-stealing in a way that minimizes memory access latency while still maintaining the flexibility and responsiveness of work-stealing. Unlike Psaroudakis et al., we do not rely on predefined thresholds to switch strategies. Instead, our scheduler uses a more fine-grained approach, dynamically adapting to the NUMA system’s topology and data locality, which is critical for maximizing performance in a variety of parallel computing scenarios. 3.3 Hierarchy Aware Stealing Drebes et al. [DHD+14] present a topology-aware scheduler, designed for NUMA architectures. The authors state that a scheduler can only be independent of the type of the machine or handle task dependencies, if it is based on the memory hierarchy of the system. This statement does not interfere with our approach of picking the locality as the highest priority, because our scheduler is specifically designed for NUMA architectures, but still is independent from its exact topology. For their topology-aware work-stealing, the authors define different parameters: a set of all processing units (PUs), an ordered set of instances of the underlying hierarchy, a function that maps for each hierarchy-instance and PU how many other PUs share the specific instance, called siblings, and another function that maps the n-th sibling of another PU on a given hierarchy level. The steal-process than can be described, by picking the n-th sibling of the same hierarchy level, with n being random, and steal from the resulting worker, if it is not the current worker, that way they prevent stealing from itself. This process can be looped, by a modifiable amount. If the steal was not successful, the same process starts over again, one level higher. The authors claim that by stepping up the hierarchy, load balancing will be achieved implicitly. They do not specify how much work gets stolen from the victims. This means that their approach suffers from the limitation that load imbalance could still occur, if the local node has still capacity available. This could result in always just working on one node, as the authors do not specify, how their workload is initially distributed across the NUMA nodes. Another problem is the amount of work that gets stolen. Stealing more than one task can improve load balancing [DLS+09], but the authors do not describe how much work actually is stolen. HotSLAW by Min et al. [MIY11] introduces another work-stealing approach that is also focused on the memory hierarchy of the underlying system. The authors adopt the help-first policy for their scheduler. This is done by using a global task queue, where the new tasks are stored. The workers then steal tasks from this queue. A problem with this approach is that the global task queue can become a bottleneck, as all workers have to access it. This can lead to contention and performance degradation. The authors claim that their work-stealing is topology aware. To achieve that they propose hierarchical victim selection and hierarchical chunk selection, in order to determine how much work gets stolen from the victim. The hierarchical victim selection is done by moving up the hierarchy, from cache, to socket, to node. The number of attempted steals per hierarchy level is modifiable. Another modifiable parameter is the chunk size. This refers to the amount of work that is stolen from the victim. With the hierarchical chunk selection, the authors propose an increase of the chunk size, when moving up the hierarchy, hence the distance to the victim. This reduces the number of remote steals, which concludes that this approach can also improve a work-stealing scheduler that follows the work-first policy. Other works [DLS+09; OP08] introduce the steal-half policy, which means that the thief always tries to steal half of the victims work. Dinan et al. [DLS+09] claim that this improves Cilks [BJK+95] approach of always stealing the oldest task. The problem with Cilk’s technique is that it only supports strict computations—that 29 3 Related Work is, computations where parent tasks only depend on the results of their child tasks after those tasks have completed. This strict, well-nested execution model enables efficient scheduling and work stealing. However, if we introduce mechanisms like coroutines, function execution can become non-strict, meaning that tasks may suspend and resume at arbitrary points and may depend on partially completed computations. This breaks the assumptions required by Cilk’s scheduler and complicates dependency tracking and load balancing. For this scenario, Dinan et al. [DLS+09] concluded that Cliks approach leads to load imbalancing. However, stealing half of the victims work leads to better load balancing and thus to a better performance. A problem, that is not taken into credit in the previous works, is the analysis of the amount of work in the specific regions. This means that work simply gets stolen, without awareness where how much work is available. The victim selection is just based on memory hierarchy or randomness. 30 4 Design and Implementation In this chapter, we present the design and implementation of our NUMA-aware cooperative scheduler. The guiding principle behind our design is to maximize the utilization of multicore systems by minimizing cross-node memory accesses and reducing synchronization overhead. This is achieved through a coroutine-based runtime model that leverages the topology information of the underlying hardware. Using this information, we implement different work-stealing mechanisms that ensures data locality, balance the workload, and reduce latency. In the following sections, we will first outline how the scheduling strategy works, which aspects we aim to optimize, and how these optimizations are achieved at a high level of abstraction. We will then proceed to a detailed explanation of the implementation. 4.1 The Runtime Environment The provided runtime environment offers two C++-classes for scheduling and executing tasks: the worker_pool and the worker. The worker corresponds to a single thread of execution, while the worker_pool is a collection of workers that manages the workers and the tasks. The worker_pool is responsible for creating workers, assigning the workers to processing units, assigning tasks to the workers, managing the amount of workers and starting and stopping the workers. The workers on the other hand execute the tasks, manage the task queues, which are decentralized across all workers, and handle their lifecycle. 4.1.1 The Worker Lifecycle Workers get created by the worker_pool’s constructor. Initially, no work is passed to them, each worker is created idle. This means that we have to implement a method that passes tasks initially to the workers, as soon as they appear. This is done in the worker_pool::run() function. This procedure is called “task scheduling” and is one of the two assignments of the scheduler. In order to design and implement our scheduler, we will later take a look at how this function can be realized in Section 4.4. After the initial work was assigned, a worker executes it’s worker::work() function. As long as the task queue of a worker is not empty, the worker executes these tasks. When the queue becomes empty, the worker will try to steal tasks from other workers. Whenever a steal-attempt is unsuccessful, the worker will go to sleep and wait for the worker_pool to wake it up again with new tasks. In order to implement the work stealing, we will adapt the worker_pool::steal() and the worker::steal_work() functions. When a worker’s queue is empty, it calls the worker_pool::steal() function to steal work. The worker_pool then selects a victim-worker, and calls the worker::steal_work() function of the victim. In this function, one or more coroutine handle(s) get taken out of the victim’s queue and put into the thief’s queue, which was passed as a reference to the two stealing-functions. This procedure can be seen in Algorithm 4.1 and is visualized in Figure 4.1. For the visualized version, 31 4 Design and Implementation Algorithm 4.1 General Stealing Algorithm procedure worker::work if own_work.size() = 0 then 𝑝𝑜𝑜𝑙.𝑠𝑡𝑒𝑎𝑙 (𝑜𝑤𝑛_𝑤𝑜𝑟𝑘) end if end procedure procedure worker_pool::steal(Queue thief_work) 𝑣𝑖𝑐𝑡𝑖𝑚 ← 𝑠𝑒𝑙𝑒𝑐𝑡_𝑣𝑖𝑐𝑡𝑖𝑚() // Done via different stealing strategies 𝑣𝑖𝑐𝑡𝑖𝑚.𝑠𝑡𝑒𝑎𝑙_𝑤𝑜𝑟𝑘 () end procedure procedure worker::steal_work(Queue thief_work) 𝑠𝑡𝑜𝑙𝑒𝑛_𝑤𝑜𝑟𝑘 ← 𝑐ℎ𝑜𝑜𝑠𝑒_𝑤𝑜𝑟𝑘 (𝑜𝑤𝑛_𝑤𝑜𝑟𝑘) // Done via different stealing strategies 𝑡ℎ𝑖𝑒 𝑓 _𝑤𝑜𝑟𝑘.𝑒𝑛𝑞𝑢𝑒𝑢𝑒(𝑠𝑡𝑜𝑙𝑒𝑛_𝑤𝑜𝑟𝑘) 𝑜𝑤𝑛_𝑤𝑜𝑟𝑘.𝑒𝑟𝑎𝑠𝑒(𝑠𝑡𝑜𝑙𝑒𝑛_𝑤𝑜𝑟𝑘) end procedure note that the worker::work() function is incomplete and only shows the scenario of the worker trying to steal work. We neglected the second worker, from whom the would be stolen, because of simplicity reasons. The effects of the steal on the workers task queues is displayed in Figure 4.2. In Figure 4.1, we only visualized the stealing process of a steal attempt of worker 1 from worker 2. Note here that the sequence diagram does not show the exact workflow of the worker, just the function calls to and from the worker_pool that are necessary for task stealing. In Section 4.3 we will take a look at different approaches for the stealing process, how they are designed and how they are implemented. Once all work has been done, the workers stay idle until the worker_pool shuts them down by calling the worker::shutdown() function. This function gracefully terminates the underlying thread. 4.1.2 Topology Information Since we aim to optimize the scheduling for NUMA systems, we need to gather topology information about the underlying hardware. This information is provided by the hwloc library [BCM+10], which is a portable library that provides a portable abstraction of hardware topology. The topology layers that can be analyzed with hwloc contain NUMA nodes, processor packages, shared caches, physical cores and processing units (also referred to, as CPUs or PUs), which represent the logical cores for hyperthreading. The main focus of using hwloc in our work is to gather information about the underlying hardware, but it is also used for thread-pinning to certain PUs, and thus to certain NUMA nodes. 4.1.3 Task Handling The tasks that can be executed runtime are based on coroutines, which were introduced to C++ in the C++20 standard [20]. The runtimes provides a full implementation of a coroutine library, called corolib. Coroutines can be used by writing classical C++-functions, but with the return- type corolib::coro, where the actual return-type of the function is passed as a generic and 32 4.1 The Runtime Environment Figure 4.1: Sequence diagram of the interacting between the worker and the worker_pool, for the stealing process of “worker 1”. return statements have to be replaced by co_return. The usage of the coroutines depends on the context from where the coroutine gets called. If the call is inside of another coroutine, i.e. an asynchronous context, coroutines should be executed by calling co_await coro(). In a synchronous context, e.g. the main function, coroutines can be called with the corolib::synchronous(coro()) function, which will ensure that the coroutine is fully executed, before continuing the con- trol flow. The awaiting of the completion of coroutines can be generalized by the usage of corolib::when_all(coro_a(), coro_b(), ...) or corolib::when_any(coro_a(), coro_b(), ...). These functions allow the execution of multiple coroutines that are started at the same time, i.e. they will get executed concurrently. This behavior follows the principle of structured concurrency, where the lifetime of spawned coroutines is lexically bound to the surrounding context. All concurrently started coroutines are joined or awaited before execution proceeds, which improves safety and 33 4 Design and Implementation Figure 4.2: Visualized effects of work stealing on the task queues. Worker 3 stole task E from worker 2. clarity in asynchronous programs [CY22]. Critical sections can be covered by the corolib::mutex. The mutex-methods lock() and unlock() are implemented as coroutines. Tasks can be handled by passing the coroutine-handle as a function parameter. The task queue of the workers is implemented as a dequeue of coroutine handles, meaning that passing these handles will become our main assignment, when implementing the work-stealing and task scheduling. When a coroutine gets called, its handle gets passed to the worker_pool and from there assigned to a worker, by enqueuing it into the worker’s queue. The worker then executes the task by calling handle.resume(). The procedure of stealing is very similar. 4.2 Baseline Scheduler Scheduling consists of two responsibilities: initial task scheduling and dynamic task stealing. Before evaluating different work-stealing strategies, we define a baseline scheduler that only handles the initial task assignment without any work stealing.This allows us to compare all stealing strategies 34 4.3 Stealing Strategies Algorithm 4.2 Round Robin Task Scheduling procedure worker_pool::run(coro_handle task) 𝑤𝑟𝑎𝑝 ← 𝑐𝑢𝑟𝑟𝑒𝑛𝑡_𝑐𝑜𝑛𝑐𝑢𝑟𝑟𝑒𝑛𝑐𝑦 // Get number of available workers 𝑖𝑛𝑑𝑒𝑥 ← ((𝑐𝑢𝑟𝑟𝑒𝑛𝑡𝑤𝑜𝑟𝑘𝑒𝑟 + 1)%𝑤𝑟𝑎𝑝) 𝑤𝑜𝑟𝑘𝑒𝑟𝑠[𝑖𝑛𝑑𝑒𝑥] .𝑒𝑛𝑞𝑢𝑒𝑢𝑒(𝑡𝑎𝑠𝑘) end procedure against a consistent, steal-free reference point. The baseline scheduler uses a round-robin strategy to distribute tasks evenly among the available workers. This simple approach ensures fairness, but does not account for load balancing, memory locality, or NUMA topology. These limitations are accepted, as the baseline scheduler only performs the initial task assignment—each task is placed once and never reassigned. Work stealing is explicitly disabled in the baseline. This isolates the impact of each stealing strategy in performance measurements, allowing for a clear comparison both against the baseline and between the strategies themselves. As described in Section 4.1.1, the initial task scheduling is located in the worker_pool::run() function, which is invoked by the runtime when new tasks arrive. The algorithm can be seen in Algorithm 4.2. A static atomic variable current determines the next worker to receive a task. It must be static to persist across function calls, and atomic to prevent race conditions when accessed from multiple threads. To avoid overloading the runtime management thread, the current thread never executes the task itself—instead, it is enqueued into the selected worker’s queue. 4.3 Stealing Strategies In the following subsections, we present the design and implementation of the three stealing strategies: NUMA-aware Stealing (Section 4.3.1), System Load-aware Stealing (Section 4.3.2) and the User-defined combination of the Stealing strategies (Section 4.3.3). The subsections will feature their design, and how they will be implemented. All three implementations will be based on the Baseline scheduler, and thus will not implement their own initial task scheduling. Another parameter in work-stealing that was discussed previously, was the amount of work that should be stolen in each stealing attempt. Initially, all three implementations will use the steal-half policy, which was introduced in Section 3.3. By using the same chunk size across the different implementations of the stealing strategies, we make sure that the chunk size is not responsible for any performance differences. This means that these differences were completely caused by the different strategies, as intended. 4.3.1 NUMA-Aware Stealing Design Our first attempt of implementing work-stealing, is to use the NUMA topology as a heuristic to select the victim worker. This heuristic can be derived by following the NUMA-aware programming principle. This concludes that we aim to steal from workers that are on the same NUMA node, as the current worker, i.e. the thief. If there is no work available on the shared NUMA node, we try 35 4 Design and Implementation Algorithm 4.3 NUMA-Aware Stealing procedure worker_pool::steal(unsigned int NUMA, Queue thief_work) 𝑤𝑜𝑟𝑘𝑒𝑟𝑠_𝑜𝑛_𝑛𝑜𝑑𝑒 ← 𝑁𝑈𝑀𝐴_𝑡𝑜_𝑤𝑜𝑟𝑘𝑒𝑟𝑠[𝑛𝑢𝑚𝑎] for all worker in workers_on_node do 𝑤𝑜𝑟𝑘𝑒𝑟.𝑠𝑡𝑒𝑎𝑙_𝑤𝑜𝑟𝑘 (𝑡ℎ𝑖𝑒 𝑓 _𝑤𝑜𝑟𝑘) if !thief_work.empty() then return end if end for 𝑠𝑜𝑟𝑡𝑒𝑑_𝑛𝑜𝑑𝑒𝑠← ℎ𝑤𝑙𝑜𝑐.𝑠𝑜𝑟𝑡_𝑛𝑜𝑑𝑒𝑠(𝑛𝑢𝑚𝑎) for all node in sorted_nodes do 𝑤𝑜𝑟𝑘𝑒𝑟𝑠_𝑜𝑛_𝑛𝑜𝑑𝑒 ← 𝑁𝑈𝑀𝐴_𝑡𝑜_𝑤𝑜𝑟𝑘𝑒𝑟𝑠[𝑛𝑜𝑑𝑒] 𝑠𝑡𝑒𝑎𝑙_𝑎𝑡𝑡𝑒𝑚𝑝𝑡𝑠() // Analogously to for-all above end for end procedure to steal from the next nearest node. The distances between the NUMA nodes are defined by the expected memory access latencies. Algorithm 4.3 shows this stealing approach. To achieve this, each worker maintains a list of all NUMA nodes, sorted by their distance to its own node. The steal process can be divided into two phases: a local steal attempt and remote steal attempts. The first phase, the local steal attempt, is performed by iterating over all workers on the same NUMA node as the thief. If a worker has work available, the thief tries to steal from it. As soon as a successful steal was performed, the function returns. If no work was available on the local NUMA node, we proceed to the second phase, the remote steal attempts. For this, we use a list of NUMA nodes that is sorted by the latency to the local node. We iterate over this list and try to steal from the workers of the current NUMA node. This is done analogously to phase one. Implementation The implementation of the NUMA-aware stealing strategy is done in the worker_pool::steal() function, with helper functions that are located in the hwlocClass class that manages the application programming interface (API) calls of the hwloc library. As we mentioned above, the stealing process is divided into two phases: the local steal at- tempt and the remote steal attempts. In the first phase, the possible victims are the work- ers that are on the same NUMA node as the thief. These workers can be retrieved from a std::unordered_map>numa_to_workers_, where all the work- ers are grouped by their NUMA node, represented by the unsigned int. This map is filled incrementally, when the workers are created. For the victim selection we just iterate over the local workers and try to steal work from them. As soon as one steal attempt was successful, we return from the function. If there was no successful steal attempt from the local workers, we proceed to the second phase, the re- mote steal attempts. Like we mentioned above, we now aim to steal from the next nearest NUMA node. For that we introduce a static, thread local std::vector>, which contains the NUMA nodes and their latencies to the local NUMA node. The unsigned int represents the NUMA node and the size_t displays the memory latency to that node. This vector is 36 4.3 Stealing Strategies chosen to be static, because it only has to be computed once per thread, since the distances between the NUMA nodes won’t change. The vector is thread local, so each worker has its own local view of the NUMA nodes and their latencies. We compute this vector by using the hwloc library, which provides a function to retrieve a distance matrix of all NUMA nodes. A helper function then extracts the remote nodes and their latencies, relative to the current NUMA node from this matrix. To iterate over the nearest NUMA nodes, we just have to sort this vector according to the latencies. The steal attempts are now done analogously to the first phase, by iterating over all workers on the current NUMA node. 4.3.2 System Load-Aware Stealing The second stealing strategy that we will look into, is the system load-aware stealing. This strategy is based on the idea of stealing from the workers where the most work is available. By doing so, we actively balance the workload from regions with a high workload to regions with a low workload. This leads to a better work distribution and thus avoiding bottlenecks, which increases the overall performance. Design The design of our system load-aware thief is based on the idea of Tzilis et al [TPTS17], which we already discussed in Section 3.2. In their work SWAS, the victim selection in the stealing process is done with decentralized information about the working statuses of all other cores of the system. These cores get grouped based on their memory hierarchy. In our implementation, we will use the same approach, but we will not use the cores themselves, but the workers. By using the workers instead of the cores, we consider a deeper level of the memory hierarchy, since the workers correspond to the threads of execution, and thus to the processing units. This concludes that our implementation will be more fine-grained, leading to better load balancing even for hyperthreading. Another major difference is the grouping of the cores, respectively the workers. We will not group the workers based on their memory hierarchy, but based on their NUMA nodes. By doing that we use a better heuristic for locality. To store the workload information decentralized, we have to store the working state of all workers in each worker. We do this by using a bitmap with two-bit entries. Each entry corresponds to a worker and contains the following information: the worker is either idle, busy or in a don’t know state, which will be our default. If a worker is idle, it is not executing any tasks and will try to steal. Busy workers are currently executing tasks, meaning that idle workers can try to steal from them. A don’t know entry means that the worker has no information about the other worker’s state. This information is gathered via stealing attempts. Tzilis et al. [TPTS17] describe how a stealing worker can collect this information. If the steal attempt is local, the thief only updates the state of the specific victim worker. The state is changed to busy, if the attempt is successful. If not, the state is updated to idle. If the steal attempt is remote, the thief updates the state of all workers in the victims’ region, i.e. the NUMA node, to busy or idle, depending on the success of the attempt. With that, the thief gathers workload information about the other regions, if no local work is available. This concludes that the don’t know entries of a worker get replaced by the “actual” information, only if the worker is trying to steal work. The next step is the victim selection. To select a victim by following the load-aware stealing strategy, 37 4 Design and Implementation we need another helping data structure. We call this data structure the NUMA estimator. In Tzilis et al.’s SWAS [TPTS17], this data structure was introduced as the region estimator. The NUMA estimator is used to create a create a score that describes the expected workloads of the NUMA nodes, based on the bitmap of the thief. This score is a weighted sum of the number of busy and don’t know workers for each NUMA node. The nodes are then sorted by their scores in a descending order. The thief tries to steal from the workers of the node with the highest score. After the stealing attempt, the thief updates it’s bitmap, as we described above. The stealing procedure can be seen in Algorithm 4.4. A flaw to this approach is that the thief will have no “actual” workload information when it steals the first time, because all entries in the bitmap are set to don’t know. This means that the thief will have equal estimators for all NUMA nodes, in which case the traversing order of the NUMA nodes doesn’t change. But this is not a problem, since by design, the local node is always the first node in the list. This means the first steal attempts will always be local. If these attempts are not successful, the next node may not be the next nearest node, nor the one with the most work, meaning that these steal attempts may result in a higher latency. The stealing algorithm is designed to steal work from these regions that are under heavy load to these regions that have less work. To do so, it uses the state of the workers to update a bitmap that is used to select the next victim. These states are updated while steal attempts are performed. The problem with this method is that the states of the workers themselves are not updated. This means that each worker has to update it’s own state in it’s lifecycle. While working, the state is set to busy. As soon as the worker has no more work, it will be set to idle. The bitmap has to be initialized, once the worker is created. Implementation Before we can implement the system load-aware stealing algorithm, we first need to implement the helping data structures. After that, we can implement the algorithm itself. The implementation of the algorithm is done in the worker_pool::steal() function, analogous to the NUMA-aware stealing. The Data Structures The first data structure we have to implement is the bitmap, which stores the state of each worker from each worker’s perspective. This bitmap is implemented as a vector of the unit8 type. This primitive data type is an unsigned integer with a size of 8 bits. This means that we can store four entries in one unit8, while each entry corresponds to one worker. To utilize the bitmap, we implement three helper functions: init_status_bitmap(), set_status(worker, status) and get_status(worker). Additionally, we define the different states: • BUSY: 0𝑏00 • IDLE: 0𝑏01 • DONTKNOW: 0𝑏10 To initialize the bitmap, all entries are set to the DONTKNOW state. Since each worker needs to maintain its own view of the global system state, the bitmap is stored per worker and includes the status of all other workers. The helper functions allow reading and writing the two-bit status entries corresponding to a specific worker. These updates are performed efficiently via compact bit 38 4.3 Stealing Strategies Algorithm 4.4 System Load-Aware Stealing procedure worker_pool::steal(unsigned int NUMA, Queue thief_work) for all [numa, workers] in numa_to_workers do // Update estimator for each node 𝑠𝑢𝑚 ← 0.0 for all worker in workers do 𝑠𝑡𝑎𝑡𝑢𝑠← 𝑡ℎ𝑖𝑒 𝑓 .𝑔𝑒𝑡_𝑠𝑡𝑎𝑡𝑢𝑠(𝑤𝑜𝑟𝑘𝑒𝑟) 𝑠𝑢𝑚+ = 𝑣𝑎𝑙𝑢𝑒_𝑜 𝑓 _𝑠𝑡𝑎𝑡𝑢𝑠() // According to BUSY, DON’T KNOW or IDLE end for 𝑒𝑠𝑡𝑖𝑚𝑎𝑡𝑜𝑟 ← 𝑠𝑢𝑚/𝑤𝑜𝑟𝑘𝑒𝑟𝑠.𝑠𝑖𝑧𝑒() 𝑡ℎ𝑖𝑒 𝑓 .𝑛𝑢𝑚𝑎_𝑒𝑠𝑡𝑖𝑚𝑎𝑡𝑜𝑟𝑠_.𝑒𝑚𝑝𝑙𝑎𝑐𝑒_𝑏𝑎𝑐𝑘 (𝑛𝑢𝑚𝑎, 𝑒𝑠𝑡𝑖𝑚𝑎𝑡𝑜𝑟) end for 𝑡ℎ𝑖𝑒 𝑓 .𝑛𝑢𝑚𝑎_𝑒𝑠𝑡𝑖𝑚𝑎𝑡𝑜𝑟𝑠_.𝑠𝑜𝑟𝑡 (𝑔𝑟𝑒𝑎𝑡𝑒𝑟, 𝑝𝑎𝑖𝑟 :: 𝑠𝑒𝑐𝑜𝑛𝑑) for all [target_numa, _] in thief.numa_estimators_ do 𝑤𝑜𝑟𝑘𝑒𝑟𝑠← 𝑛𝑢𝑚𝑎_𝑡𝑜_𝑤𝑜𝑟𝑘𝑒𝑟𝑠[𝑡𝑎𝑟𝑔𝑒𝑡_𝑛𝑢𝑚𝑎] for all victim in workers do 𝑣𝑖𝑐𝑡𝑖𝑚.𝑠𝑡𝑒𝑎𝑙_𝑤𝑜𝑟𝑘 (𝑤𝑜𝑟𝑘) if 𝑛𝑢𝑚𝑎 == 𝑡𝑎𝑟𝑔𝑒𝑡_𝑛𝑢𝑚𝑎 then // Local Steal Update 𝑡ℎ𝑖𝑒 𝑓 .𝑠𝑒𝑡𝑆𝑡𝑎𝑡𝑢𝑠(victim, steal_success ? BUSY : IDLE) else // Remote Steal Update for all worker in numa_to_workers[target_numa] do 𝑡ℎ𝑖𝑒 𝑓 .𝑠𝑒𝑡𝑆𝑡𝑎𝑡𝑢𝑠(worker, steal_success ? BUSY : IDLE) end for end if end for end for end procedure manipulations, allowing fast access without requiring additional data structures. The second data structure is the NUMA estimator, implemented as a vector of pairs, where each pair contains a NUMA node ID and a corresponding load estimate. This estimator is maintained per worker and can be used to track observed or inferred load on different NUMA nodes. It enables workers to make informed decisions during work stealing, such as prioritizing nodes with higher estimated workload. Both data structures are publicly accessible to support decentralized coordination. Steal function The stealing algorithm is, as mentioned above, implemented in the worker_pool::steal() function. The algorithm can be divided into three parts: • The Update of the estimators • The Victim Selection + Stealing • The Update of the bitmap First, we have to update the estimators. For that, the NUMA estimators vector of the thief has to be brought up to date. As the NUMA estimator stores the expected workload of each NUMA node, we have to iterate over all NUMA nodes and calculate the expected workloads for each node. The 39 4 Design and Implementation iteration over the NUMA nodes is done via the numa_to_workers_ map. The update of the workload for each NUMA node is done by accessing each of its worker’s bitmaps and extract the states of the workers. If a worker is busy, we add 1.0 to the expected workload, if it is in the don’t know state, we add 0.5. The idle state can be neglected, since we would add 0.0. Each load estimation, i.e. this sum in one iteration over the NUMA nodes, is normalized and then stored in the NUMA estimator of the thief. Secondly, for the victim selection, we now want to extract the NUMA node with the highest expected workload. This is done by sorting the NUMA estimators vector in a descending order, based on the expected workloads of each NUMA node. By iterating over the sorted vector, we implicitly try to steal from the NUMA nodes with the highest expected workloads first. Like in remote steal attempts of the NUMA-Aware stealing in Section 4.3.1, we try to steal from the workers of the current NUMA node by iterating over them. For each steal attempt, we have to store the result, if it was successful or not. This information is needed for step three, the updating of the bitmap. For that, we have to check if the steal was local or remote. In the local case, we update the state of the victim accordingly to the success of the steal attempt. In the remote case, we update the state of all the workers of the victim’s NUMA node, again accordingly to the success of the steal attempt. As soon as we perform a successful steal attempt, we return from the function, analogously to the NUMA-Aware stealing. State Handling At last, we need to adopt the worker::work() function that defines the lifecycle of the workers. If the worker has work, the state in the own bitmap is set to busy. If the worker has no work, i.e. before it tries to steal, the state is set to idle. If the steal is successful, the worker get’s back into the work() loop and the state is set to busy again. If the steal was not successful, the worker goes to sleep and the state stays idle. When the worker is woken up by the worker_pool, the state is set to busy again, as soon as the worker starts executing tasks. Before all of that, the init_status_bitmap() function is used in the constructor of the worker. 4.3.3 Combining NUMA- and System Load-Aware Stealing The victim selection of the NUMA-aware stealing strategy is purely based on the NUMA locality. Thus, the victim selection is only based on making assumptions and measurements about the memory latencies between the different NUMA nodes. This means that a worker would still steal from its local node, even if there was only one available task, while for example on the next nearest node, all workers are fully utilized. This problem is addressed by the system load-aware stealing strategy, where a steal is always attempted from the NUMA node with the highest available workload. The flaw of the load aware stealing strategy is that it makes no assumptions about the NUMA locality, leading to potentially more remote stealing and thus to increasing memory latencies. Thus, we present a third stealing approach, NUMA-locality and Load-aware Stealing Strategy (NLLASS) that combines the benefits of both previously mentioned strategies. 40 4.3 Stealing Strategies Design In our combined approach, we introduce a weighted evaluation of potential victim NUMA nodes that considers both estimated system load and NUMA distance. The goal is to make a balanced stealing decision that accounts for memory locality while still targeting regions with the highest expected workload. The core idea is to compute a scheduling score for each NUMA node that determines the priority order when attempting to steal work. First, the load estimators are computed for each NUMA node. These estimators are the same that we described in the system load-aware stealing strategy in Section 4.3.2, with the minor distinction that we use one minus their original values. This means that a smaller estimator indicates the higher workload. Secondly, each NUMA node, except the local one, is assigned a scheduling score by multiplying its load estimator with its distance from the current worker’s NUMA node. These distances are the same we described in Section 4.3.1. By using these NLLASS load estimators, this leads to an artificial decrease of the latency, far-away nodes with a high workload appear closer. This ensures that NUMA nodes that are both nearby and expected to have available work receive a higher priority. The victim selection is based on this scheduling score, where a lower score indicates a better suited NUMA node to steal from, than a higher score. This hybrid design ensures locality-preserving stealing when possible, but also leverages system-wide information to avoid underutilized remote cores. The local NUMA node is always included with a fixed distance score of zero to allow for immediate local stealing if possible. Another possibility to gain the perfect score zero is, if all workers on a node are busy, leading to an estimator of zero. This means that fully utilized NUMA nodes are treated the same way as the local node. This could become a problem, because if the system is nearly fully utilized, it could lead to a massive increase in remote steals. We tackle this problem the same way as in Section 4.3.2, by placing the local node always at the top of the queue. This means that the workers of the local node are always tried before workers of remote nodes, if the scheduling score is equal. A simplified version of the algorithm can be seen in Algorithm 4.5. Example. We consider 3 NUMA nodes, each with 3 workers. The table below shows the state of all workers and their distance to node 0. We display the victim selection process for Worker 3 on node 0: Node Worker 1 Worker 2 Worker 3 Distance to Node 0 0 IDLE IDLE IDLE 0 1 BUSY IDLE IDLE 2 2 BUSY BUSY IDLE 3 Table 4.1: Example of the workers’ states and their distances from node 0, the thief’s node. We compute the workload estimators and scheduling scores for node 0 as follows: 41 4 Design and Implementation estimator0(1) = 1 − 1 + 0 + 0 3 = 0.66 estimator0(2) = 1 − 1 + 1 + 0 3 = 0.33 sched_score0(1) = 0.66 · 2 = 1.32 sched_score0(2) = 0.33 · 3 = 0.99 Since node 2 has the lowest scheduling score, Worker 3 from node 0 would choose to steal from Worker 1 or 2 on node 2. Note, that we ignore the local node in this example, to demonstrate the effects of the scheduling score. In this scenario, stealing from the local node is not possible, as all workers on this node are idle. Implementation The NUMA-aware and the load-aware stealing strategies rely on information about the system workload and the memory latencies between the NUMA nodes, respectively. To combine the two stealing strategies, we need both pieces of information. Thus, we introduce a new data structure to store the combination the two strategies. After collecting all the necessary data, we combine the two stealing algorithms. Data structures To analyze the systems workload, we use the decentralized storage of the workers states in each worker via a bitmap (Section 4.3.2). In the system load aware stealing strategy we used a std::vector>, called numa_estimators_ to store the estimated workloads of the NUMA nodes. For the combined stealing we change this data structure into a std::unordered_map. This simplifies the access to the estimator values for specific nodes. Since we don’t have to sort the NLLASS load estimators, this change won’t cause any problems. Algorithm 4.5 Combining Stealing Strategies procedure worker_pool::steal(unsigned int numa, Queue thief_work) 𝑡ℎ𝑖𝑒 𝑓 .𝑛𝑢𝑚𝑎_𝑒𝑠𝑡𝑖𝑚𝑎𝑡𝑜𝑟𝑠_← 𝑐𝑎𝑙𝑐𝐸𝑠𝑡𝑖𝑚𝑎𝑡𝑜𝑟𝑠() for all [target_numa, distance] in current.sorted_numa_nodes do 𝑡ℎ𝑖𝑒 𝑓 .𝑠𝑐ℎ𝑒𝑑_𝑠𝑐𝑜𝑟𝑒_← 𝑐𝑎𝑙𝑐𝑆𝑐𝑜𝑟𝑒(𝑡ℎ𝑖𝑒 𝑓 _𝑒𝑠𝑡𝑖𝑚𝑎𝑡𝑜𝑟, 𝑑𝑖𝑠𝑡𝑎𝑛𝑐𝑒) end for 𝑡ℎ𝑖𝑒 𝑓 .𝑠𝑐ℎ𝑒𝑑_𝑠𝑐𝑜𝑟𝑒_.𝑠𝑜𝑟𝑡 (𝑙𝑒𝑠𝑠, 𝑝𝑎𝑖𝑟 :: 𝑠𝑒𝑐𝑜𝑛𝑑) for all [target_numa, _] in current.sched_score_ do 𝑡𝑟𝑦_𝑡𝑜_𝑠𝑡𝑒𝑎𝑙 (𝑡𝑎𝑟𝑔𝑒𝑡_𝑛𝑢𝑚𝑎.𝑤𝑜𝑟𝑘𝑒𝑟𝑠()) 𝑢𝑝𝑑𝑎𝑡𝑒𝐵𝑖𝑡𝑚𝑎𝑝(𝑡𝑎𝑟𝑔𝑒𝑡_𝑛𝑢𝑚𝑎.𝑤𝑜𝑟𝑘𝑒𝑟𝑠()) end for end procedure 42 4.4 Advanced Task Scheduling - NLHScheduler For the NUMA distances, we introduce a data structure, implemented as a std::vector>. In this new data structure, called sorted_numa_nodes_, we store all the distances to each NUMA node, relative to the worker’s local node. This vector is computed the same way as the sorted NUMA nodes vector in the NUMA aware stealing strategy. Additionally, in a std::vector> we store the scheduling infor- mation that takes NUMA-aware and load-aware information into account. This new data structure is called “sched_scores”. The unsigned int denotes the NUMA node, while the double represents the scheduling score associated with that node. Steal function The stealing algorithm is located in the worker_pool::steal() function. The algorithm is divided into three phases: 1. Collecting the scheduling data 2. Selecting the victim 3. Updating the system load data The collection of the scheduling data also has two parts: First, we have to collect the system load information and then secondly combine it with the NUMA distance information. The collection of the system load information is analogous to the updating of the estimators in Section 4.3.2, with two small differences. On the one hand, the handling of the estimators changes, because we changed the numa_estimators_ from a vector to a map. On the other hand, the estimator value is stored differently. Instead of storing the estimator directly, as done in the system load-aware strategy, we store its complement, as previously described. To compute the scheduling score, we iterate over the sorted NUMA nodes vector of the thief. In each iteration, we add a new entry to the sched scores vector. This entry is the current threads NUMA node and the multiplication of the current latency with the current load estimator. Note, that the sched_scores vector should be cleared before this update, to avoid data inconsistencies. The local node has to be inserted manually, since it is not present in the sorted NUMA nodes vector. These scores are now sorted by the scheduling scores. After that we start with phase two and three of the algorithm. These two phases are almost identical to phases two and three of Algorithm 4.4. The small difference is that, instead of iterating over the numa_estimators_, we now iterate over the sched_scores vector. 4.4 Advanced Task Scheduling - NLHScheduler After implementing and evaluating different stealing strategies, we now turn our attention to the initial task scheduling. We discuss how to improve the first assignment of tasks to the workers. In the previous stealing strategy implementations this was done by using a round-robin approach. However, round-robin scheduling is not always the best approach, since it does not consider NUMA locality or the current system load distribution. We implement the task scheduling algorithm on top of the combined stealing strategy, creating the NLHScheduler. In this section, we will discuss how we improve the initial task scheduling by combining the NUMA locality and the system load distribution as parameters for the scheduling decision. First, we will discuss the design of the new scheduling strategy. After that, we will take a look at the implementation. 43 4 Design and Implementation Algorithm 4.6 Advanced Task Scheduling procedure worker_pool::run(coro_handle task) 𝑙𝑜𝑐𝑎𝑙_𝑛𝑜𝑑𝑒 ← 𝑔𝑒𝑡_𝑙𝑜𝑐𝑎𝑙_𝑛𝑜𝑑𝑒() if 𝑙𝑜𝑐𝑎𝑙_𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑 < 𝑚𝑒𝑑𝑖𝑎𝑛_𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑 then 𝑤𝑜𝑟𝑘𝑒𝑟 ← 𝑙𝑜𝑤𝑒𝑠𝑡_𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑_𝑤𝑜𝑟𝑘𝑒𝑟 (𝑙𝑜𝑐𝑎𝑙_𝑛𝑜𝑑𝑒) 𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑 [𝑙𝑜𝑐𝑎𝑙_𝑛𝑜𝑑𝑒]+ = 1 else 𝑙𝑜𝑤𝑒𝑠𝑡_𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑_𝑛𝑜𝑑𝑒 ← 𝑔𝑒𝑡_𝑙𝑜𝑤𝑒𝑠𝑡_𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑_𝑛𝑜𝑑𝑒() 𝑤𝑜𝑟𝑘𝑒𝑟 ← 𝑙𝑜𝑤𝑒𝑠𝑡_𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑_𝑤𝑜𝑟𝑘𝑒𝑟 (𝑙𝑜𝑤𝑒𝑠𝑡_𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑_𝑛𝑜𝑑𝑒) 𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑 [𝑙𝑜𝑤𝑒𝑠𝑡_𝑤𝑜𝑟𝑘𝑙𝑜𝑎𝑑_𝑛𝑜𝑑𝑒]+ = 1 end if 𝑤𝑜𝑟𝑘𝑒𝑟.𝑒𝑛𝑞𝑢𝑒𝑢𝑒(𝑡𝑎𝑠𝑘) end procedure 4.4.1 Design To update the task scheduling from the baseline round-robin approach, we consider both NUMA locality and system load distribution when assigning tasks initially to workers. The goal is to create a more efficient and balanced task distribution that minimizes memory access latencies and maximizes utilization of the available resources. To achieve this, we introduce a new data structure that stores an estimation of the workloads of the NUMA nodes. Based on the values in this data structure, we can assign tasks to the workers. To introduce NUMA locality into the scheduling decision, we also need to gather information about the NUMA node of the current thread that assigns the tasks. First, we compare the workload of the local NUMA node to the median workload of all other NUMA nodes. If the local workload is below this median, the task is assigned to a worker on the local node. However, if the local workload exceeds the median, the task is assigned to a worker on the NUMA node with the lowest workload. Secondly, once the target NUMA node is selected, we identify the specific worker on that node to execute the task. This is done by choosing the worker with the lowest workload, estimated using the size of their task queue. It is important to note that using queue size as a workload heuristic can be problematic, as it does not always accurately represent the worker’s actual processing load. If one worker has a large queue of very small tasks and another worker has a small queue of very large tasks, we would assign a new task to the worker with the smaller queue, even though it might take longer to execute the few large tasks than the many small ones. To optimize this metric, one would have to use a more sophisticated heuristic that predicts the execution times of the tasks in the queues. Since this is a difficult task, we will use the queue size as the heuristic and leave the optimization of the heuristic for future work. After selecting the worker, we assign the new task to it and update the workload information of the worker’s NUMA node. The task scheduling algorithm can be seen in Algorithm 4.6. 44 4.4 Advanced Task Scheduling - NLHScheduler 4.4.2 Implementation Like the baseline scheduler, the implementation of the initial task scheduling will be done in the worker_pool::run() function. While implementing the new scheduling strategy, we encountered some problems with the runtime environment. In this section, we will firstly discuss the implementation of the scheduling strategy and secondly take a look at the problems and how we solved them. The Run Function The implementation of the new task scheduling strategy can be found in the worker_pool::run() function that takes one parameter, which corresponds to a std::coroutine_handle<>. This parameter is the task that should be assigned. For the first step of the algorithm, we need to compare the local workload with the median workload of the other NUMA nodes. To do that, we first have to extract the local NUMA node. But there is a problem with the function that we use to determine the local NUMA node. Since we are creating the workers ourselves and thus create the thread_index of the worker ourselves, this function only works for threads that are created as workers from our runtime. The flaw of this is that the executing thread of the worker_pool::run() function is not always a worker. Here we have to differ between external and internal scheduling. If the executing thread is indeed a worker, the internal scheduling can go on and we can use the function to extract the local NUMA node. If the executing thread is not a worker, it means that we have to consider external scheduling. That means that a non-worker thread, e.g. an operating system thread, determines the scheduling decision. In this case we can’t determine the local NUMA node the way we did before and use NUMA node zero as a default. This differentiation can be done by comparing the thread index with a constant in the worker class that indicates, when a thread is not a worker. After selecting our local NUMA node, we can extract the workload estimations of the local node and the median of the other NUMA nodes. For that we use a std::unordered_map, called numa_workloads to store the workload estimations for each NUMA node. The key of the map refers to the NUMA node and the values to the estimated workloads. The median of the other NUMA nodes is calculated by using a helper function. If the local workload is smaller than the median, we extract the workers of the local node. We then filter the workers to retrieve the worker with the shortest task queue. Now we can enqueue the passed task to this worker. After that, we update the workload of the local node by adding one to the value in the numa_workloads map. For the second case, we have to extract the NUMA node with the lowest workload. This is done by filtering the numa_workloads map, After extracting the NUMA node with the lowest workload, we can select the target worker of this node the exact same way as we did it for the local node. The only difference is the update of the workload. Here, we add one to the entry of the numa_workloads map that corresponds to the NUMA node of the selected worker. 45 5 Evaluation In this chapter, we evaluate the various schedulers implemented as part of this work. We begin by describing the benchmark used to assess scheduler performance, including its design and implementation. This is followed by an explanation of the performance metrics employed in the evaluation. We then outline the parameters that are varied during the benchmark runs. Prior to presenting the results, we describe the test environment in which the experiments were conducted. Finally, we present and analyze the results of the benchmark executions. 5.1 Benchmark In the following two subsections we discuss the design and implementation of the benchmark that we use to evaluate the performance of the different scheduling strategies that were described and implemented in Chapter 4. Firstly, we take a look at thew design of the benchmark and its core features. Secondly, we discuss how the benchmark leads to work stealing and how we can create different scenarios to test the different implementations. Finally, we explain how we achieved implemented the benchmark. 5.1.1 Design To test and evaluate the performance of a NUMA-aware cooperative work-stealing scheduler, we have to create a workload that simulates realistic contention between threads. Our benchmark is designed to provoke both local and remote work-stealing events under controlled conditions. The workload consists of coroutine-based matrix operations—specifically, the generation and multiplication of large matrices. The benchmark includes the following coroutine functions, which together simulate computational load and scheduling challenges: • Matrix Generation: This coroutine generates a matrix of a specified size, fills it with randomized values, and stores it in a globally accessible container for use by other coroutines. In the benchmark, only two matrices are created and reused throughout the workload. • Matrix Multiplication: This coroutine performs the core computation by multiplying two matrices. • Multiple Multiplications: A coroutine that performs 100 consecutive matrix multiplications to generate prolonged load and occupy a worker thread for an extended period. • Thief Function: This coroutine simply awaits the completion of a multiple_multiplications coroutine. Its purpose is to simulate idle workers by blocking until another coroutine completes, thereby encouraging work stealing. 47 5 Evaluation • Start Work: This coroutine coordinates the workload. It first launches the matrix generation coroutines and waits for them to complete. Then, it starts a set of matrix multiplication and thief coroutines in an alternating fashion to mix active and idle workloads across workers. • Main Function: This non-coroutine function configures and initiates the benchmark. It configures parameters including the number of matrix multiplications, the number of thief coroutines, and the concurrency level, defined as the number of worker threads. We measure the execution time, as well as local and remote steals. Each execution of the benchmark is configured via command-line arguments, which specify the number of worker threads, the number of single matrix multiplication coroutines, and the number of thief coroutines. The benchmark is repeated multiple times—by default, five runs—to compute median values for the metrics. The outputs of the benchmark are the medians of the execution times, the performed local and remote steals and the throughput. These outputs are written into an output file. Imbalanced Workload When executed with the baseline round-robin scheduler, the benchmark distributes tasks evenly across NUMA nodes. However, this uniform distribution does not reflect realistic conditions in practical systems, where load imbalances are common. To simulate such an imbalance, we modify the matrix multiplication coroutine: if it is executed on a specific NUMA node, it performs three matrix multiplications instead of one. This artificially increases the workload on that NUMA node, creating an imbalance that encourages task stealing from workers on this node. This scenario allows us to evaluate the effectiveness of the scheduler’s NUMA-aware work-stealing strategy under asymmetric load conditions. 5.1.2 Implementation The benchmark implementation is based on the coroutine framework corolib, which allows asynchronous task definitions via coroutine syntax. The core idea is to create varying loads across worker threads and NUMA domains to provoke stealing events. Below we describe the key components in detail: Matrix Representation A Matrix struct is defined to hold square matrices using nested std::vector containers. The global container dataVec, implemented as std::vector, stores the matrices generated during the benchmark, and is protected by a coroutine-aware corolib::mutex. Matrix Generation The coroutine generate_data_coro creates a matrix of size 500 × 500, fills it with random values, and stores it in the global matrix vector. The coroutine uses co_await to acquire and release a lock when inserting data, ensuring thread-safe access. An important note here is that all generated and calculated matrices would be neglected by the compiler, due to the fact that they are not used in any further calculations. To prevent this, we use the corolib::benchmark::do_not_optimize() function. This way, we ensure that all matrix 48 5.2 Metrics generations and multiplications are actually performed and not optimized by the compiler. This function is used in the generate_data_coro, multiply_matrix_coro, and multiple_multiplications coroutines on the result matrices. Matrix Multiplication The coroutine multiply_matrix_coromultiplies two matrices and consumes significant CPU time. Each invocation logs which worker executes the coroutine. These logs are not thread-safe and thus are not printed into the output file. Their purpose was to help debugging certain implementations. After the logs, the coroutine performs one or three matrix multiplications depending on the workload configuration. This means that, based on the NUMA node of the worker, determined via the hwlocClass, additional multiplications can be performed to simulate imbalance. In the asynchronous workload configuration, workers on NUMA node zero perform the three multiplications, while workers on the other NUMA nodes perform only one. Heavy Computation Task The coroutine multiple_multiplications executes 100 matrix multi- plications in a loop. This simulates a long-running, compute-heavy task that occupies a worker thread and increases the likelihood of stealing from other workers. Thief Coroutine The coroutine start_stealer simulates idle