ScyllaDB & Seastar Architecture: Thread-per-Core Execution, Shared-Nothing C++ & P99 Latency

Table of Contents(10 sections)
ScyllaDB is a high-performance, NoSQL distributed database compatible with Apache Cassandra and Amazon DynamoDB APIs. Its architectural distinctiveness stems from its underlying C++ asynchronous engine, Seastar. This document dissects ScyllaDB's core architectural tenets: thread-per-core execution, shared-nothing design, NUMA awareness, and user-space I/O scheduling, elucidating how these principles collectively deliver predictable sub-millisecond P99 latencies even under extreme loads.
The Seastar Engine: A Shared-Nothing, Thread-per-Core Model
Traditional database systems, including Apache Cassandra, often rely on a multi-threaded architecture where a pool of threads processes requests. These threads contend for shared resources (locks, caches) and are subject to operating system scheduler interventions, leading to context switching overheads and unpredictable latency spikes. Furthermore, JVM-based systems like Cassandra introduce garbage collection (GC) pauses, which can significantly impact P99 latency.
Seastar, the asynchronous C++ framework powering ScyllaDB, fundamentally re-architects this model. It adopts a "shared-nothing, thread-per-core" execution paradigm.
Thread-per-Core Execution
In Seastar, each CPU core is assigned a dedicated Seastar thread (also known as a "shard" or "reactor"). This thread is pinned to its core and runs an event loop. All data and logic processed by that core are local to it. There is no shared memory or shared data structures between cores that require locks or atomic operations for concurrent access.
This design eliminates:
- Context Switching Overhead: The OS scheduler is largely bypassed for application-level scheduling. Each core's thread runs continuously, processing events from its local queue.
- Lock Contention: With no shared mutable state, mutexes, semaphores, and other synchronization primitives are largely unnecessary for core-to-core data access, simplifying concurrency and boosting throughput.
- Cache Invalidation: Data locality is maximized. A core primarily operates on data residing in its own L1/L2/L3 caches, reducing cache misses and inter-core cache coherence traffic.
Shared-Nothing Architecture
The "shared-nothing" principle extends beyond CPU cores to other resources. Each Seastar thread manages its own:
- Memory Pool: NUMA-aware memory allocation ensures that memory accessed by a core is allocated from the NUMA node local to that core.
- Network Stack: Seastar can bypass the kernel's network stack using technologies like DPDK (Data Plane Development Kit) or XDP (eXpress Data Path) for direct access to network interface cards (NICs). This allows each core to handle its own network I/O, reducing kernel overhead and improving packet processing rates.
- Disk I/O Queue: Each core manages its own queue of disk I/O requests, leveraging asynchronous I/O mechanisms like Linux AIO or
io_uring.
When a request arrives, it's typically hashed to a specific core based on a partition key. That core then handles the request end-to-end, from network ingress to disk I/O and back, without involving other cores unless data needs to be explicitly transferred (e.g., for cross-shard queries, which are handled via message passing).
Asynchronous Programming Model
Seastar employs a future-based asynchronous programming model. Operations that would traditionally block (e.g., disk I/O, network I/O) return a future<T> object immediately. The Seastar event loop then schedules other ready tasks while the I/O operation completes. Once the I/O is done, the corresponding future is "fulfilled," and its continuation (a callback or chained operation) is scheduled for execution on the same core.
This cooperative multitasking within a single thread-per-core avoids the overhead of OS context switches while still allowing high concurrency.
// Example: Seastar asynchronous I/O
#include <seastar/core/app-template.hh>
#include <seastar/core/future.hh>
#include <seastar/core/file.hh>
#include <seastar/core/reactor.hh>
#include <seastar/core/thread.hh> // For seastar::thread
// Function to write data asynchronously to a file
seastar::future<> write_to_file(const seastar::sstring& filename, const seastar::sstring& data) {
// Open the file asynchronously. O_CREAT | O_TRUNC | O_WRONLY are standard flags.
// 0644 is file permissions.
return seastar::open_file_dma(filename, seastar::open_flags::rw_create | seastar::open_flags::truncate).then([data](seastar::file f) {
// Allocate a DMA-aligned buffer for efficient I/O
auto buffer = seastar::temporary_buffer<char>::aligned(4096, data.size());
std::copy(data.begin(), data.end(), buffer.begin());
// Write the buffer to the file asynchronously
return f.dma_write(buffer.get(), 0, buffer.size()).then([f = std::move(f), buffer = std::move(buffer)] (size_t bytes_written) {
std::cout << "Wrote " << bytes_written << " bytes to " << f.get_path() << std::endl;
// Close the file asynchronously
return f.close();
});
}).handle_exception([](std::exception_ptr ep) {
// Handle any exceptions during file operations
std::cerr << "Error writing to file: " << seastar::current_exception_better_what(ep) << std::endl;
return seastar::make_exception_future<>(ep);
});
}
// Function to read data asynchronously from a file
seastar::future<seastar::sstring> read_from_file(const seastar::sstring& filename) {
return seastar::open_file_dma(filename, seastar::open_flags::ro).then([filename](seastar::file f) {
// Get file size to allocate buffer
return f.size().then([f = std::move(f), filename](uint64_t size) mutable {
auto buffer = seastar::temporary_buffer<char>::aligned(4096, size);
// Read into the buffer asynchronously
return f.dma_read(buffer.get(), 0, size).then([f = std::move(f), buffer = std::move(buffer)](size_t bytes_read) mutable {
std::cout << "Read " << bytes_read << " bytes from " << f.get_path() << std::endl;
// Close the file asynchronously
return f.close().then([buffer = std::move(buffer)]() mutable {
return seastar::sstring(buffer.begin(), buffer.size());
});
});
});
}).handle_exception([](std::exception_ptr ep) {
std::cerr << "Error reading from file: " << seastar::current_exception_better_what(ep) << std::endl;
return seastar::make_exception_future<seastar::sstring>(ep);
});
}
int main(int argc, char** argv) {
seastar::app_template app;
// Define the application's main function
app.run(argc, argv, [] {
return seastar::make_ready_future().then([] {
seastar::sstring test_data = "Hello, Seastar asynchronous I/O!";
seastar::sstring filename = "test_async_io.txt";
// Chain asynchronous operations: write, then read, then print
return write_to_file(filename, test_data).then([filename] {
return read_from_file(filename);
}).then([](seastar::sstring content) {
std::cout << "File content: " << content << std::endl;
}).handle_exception([](std::exception_ptr ep) {
std::cerr << "Application failed: " << seastar::current_exception_better_what(ep) << std::endl;
return seastar::make_exception_future<>(ep);
});
});
});
return 0;
}
To compile and run this Seastar example:
# Assuming Seastar is installed and SEASTAR_HOME is set
g++ -std=c++17 -Wall -Werror -O2 -I${SEASTAR_HOME}/include -L${SEASTAR_HOME}/build/lib -Wl,-rpath=${SEASTAR_HOME}/build/lib -o async_io async_io.cpp -lseastar -lfmt -lstdc++fs -lboost_program_options -lboost_thread -lboost_system -lboost_filesystem -lboost_chrono -lboost_context -lboost_atomic -lhwloc -latomic -lrt -lm -ldl -lucontext -lnuma -lz -lcryptopp -lgnutls -lprotobuf -ljsoncpp -lcap -luring
./async_io --smp 1 # Run with 1 core
NUMA Awareness
Modern multi-socket servers employ Non-Uniform Memory Access (NUMA) architectures. Memory access times vary depending on whether the memory is local to the CPU accessing it or located on another NUMA node. ScyllaDB is explicitly NUMA-aware. It partitions its memory pools such that each Seastar thread allocates memory from the NUMA node it resides on. This significantly reduces cross-NUMA node memory access, which is slower and consumes inter-socket bandwidth.
User-Space I/O Scheduling (Linux AIO / io_uring)
ScyllaDB bypasses the kernel's block layer for disk I/O whenever possible. It uses asynchronous I/O interfaces like Linux AIO (Asynchronous I/O) or, more recently and preferably, io_uring. io_uring is a modern Linux kernel interface that provides a highly efficient, zero-copy asynchronous I/O mechanism.
Each Seastar core maintains its own io_uring submission and completion queues. This allows direct submission of I/O requests from user-space to the kernel, and direct retrieval of completion events, minimizing context switches and system call overheads. ScyllaDB also implements its own I/O scheduler in user-space, allowing fine-grained control over I/O prioritization and fairness, which is crucial for maintaining low latency under mixed workloads.
// Conceptual illustration of io_uring usage within Seastar (simplified)
// In reality, Seastar abstracts this heavily.
#include <liburing.h>
#include <fcntl.h>
#include <unistd.h>
#include <sys/stat.h>
#include <iostream>
#include <vector>
#include <string>
// This is a simplified, standalone example to demonstrate io_uring concepts.
// Seastar integrates io_uring much more deeply into its reactor model.
const int QUEUE_DEPTH = 64;
const int BLOCK_SIZE = 4096;
int main() {
struct io_uring ring;
int ret = io_uring_queue_init(QUEUE_DEPTH, &ring, 0);
if (ret < 0) {
std::cerr << "io_uring_queue_init: " << strerror(-ret) << std::endl;
return 1;
}
int fd = open("test_io_uring.txt", O_RDWR | O_CREAT | O_TRUNC, 0644);
if (fd < 0) {
std::cerr << "open: " << strerror(errno) << std::endl;
io_uring_queue_exit(&ring);
return 1;
}
std::string write_data = "Hello from io_uring!";
std::vector<char> write_buffer(BLOCK_SIZE);
std::copy(write_data.begin(), write_data.end(), write_buffer.begin());
// Prepare a write request
struct io_uring_sqe *sqe = io_uring_get_sqe(&ring);
if (!sqe) {
std::cerr << "io_uring_get_sqe failed" << std::endl;
close(fd);
io_uring_queue_exit(&ring);
return 1;
}
io_uring_prep_write(sqe, fd, write_buffer.data(), write_buffer.size(), 0);
sqe->user_data = 1; // Unique identifier for this request
// Submit the request
io_uring_submit(&ring);
// Wait for completion
struct io_uring_cqe *cqe;
ret = io_uring_wait_cqe(&ring, &cqe);
if (ret < 0) {
std::cerr << "io_uring_wait_cqe: " << strerror(-ret) << std::endl;
close(fd);
io_uring_queue_exit(&ring);
return 1;
}
if (cqe->res < 0) {
std::cerr << "Write failed: " << strerror(-cqe->res) << std::endl;
} else {
std::cout << "Write completed: " << cqe->res << " bytes" << std::endl;
}
io_uring_cqe_seen(&ring, cqe); // Mark completion as seen
// Prepare a read request
std::vector<char> read_buffer(BLOCK_SIZE);
sqe = io_uring_get_sqe(&ring);
if (!sqe) {
std::cerr << "io_uring_get_sqe failed for read" << std::endl;
close(fd);
io_uring_queue_exit(&ring);
return 1;
}
io_uring_prep_read(sqe, fd, read_buffer.data(), read_buffer.size(), 0);
sqe->user_data = 2; // Unique identifier for this request
// Submit the request
io_uring_submit(&ring);
// Wait for completion
ret = io_uring_wait_cqe(&ring, &cqe);
if (ret < 0) {
std::cerr << "io_uring_wait_cqe for read: " << strerror(-ret) << std::endl;
close(fd);
io_uring_queue_exit(&ring);
return 1;
}
if (cqe->res < 0) {
std::cerr << "Read failed: " << strerror(-cqe->res) << std::endl;
} else {
std::cout << "Read completed: " << cqe->res << " bytes. Content: " << std::string(read_buffer.data(), cqe->res) << std::endl;
}
io_uring_cqe_seen(&ring, cqe);
close(fd);
io_uring_queue_exit(&ring);
return 0;
}
To compile and run this io_uring example:
g++ -std=c++17 -Wall -Werror -O2 -o io_uring_example io_uring_example.cpp -luring
./io_uring_example
Note: io_uring requires a recent Linux kernel (5.1 or newer for basic functionality, 5.8+ for advanced features).
ScyllaDB vs. Apache Cassandra: Architectural Comparison
| Feature | ScyllaDB (Seastar) | Apache Cassandra (JVM) |
|---|---|---|
| Execution Model | Thread-per-core, shared-nothing | Multi-threaded, shared-memory |
| Concurrency | Cooperative multitasking (futures) | Preemptive multitasking (OS threads) |
| Language | C++ | Java |
| Memory Management | Manual (NUMA-aware allocators) | JVM GC (Stop-the-world/concurrent) |
| I/O Model | User-space AIO/io_uring, direct | Kernel-buffered AIO/NIO, OS scheduler |
| Network Stack | Optional kernel bypass (DPDK/XDP) | Standard kernel TCP/IP stack |
| Latency Predictability | High (sub-millisecond P99) | Lower (GC pauses, context switches) |
| Resource Utilization | Near 100% CPU, high I/O efficiency | Variable, GC overhead, context switching |
| Startup Time | Fast | Slower (JVM JIT, class loading) |
| Footprint | Smaller | Larger (JVM runtime) |
Achieving Predictable P99 Latency
ScyllaDB's architectural choices directly contribute to its ability to sustain predictable sub-millisecond P99 latencies under high throughput (e.g., 500,000 writes/sec).
- Elimination of GC Pauses: Being written in C++, ScyllaDB avoids the unpredictable "stop-the-world" pauses inherent in JVM garbage collection. Memory management is deterministic and controlled.
- Reduced Context Switching: Thread-per-core pinning and the cooperative multitasking model drastically reduce OS-level context switches, which are a major source of latency jitter.
- Data Locality and Cache Efficiency: The shared-nothing design ensures data and code are local to the processing core, maximizing CPU cache hit rates and minimizing expensive memory accesses.
- Efficient I/O: User-space I/O scheduling with
io_uringprovides direct, low-latency access to storage devices, bypassing kernel overheads and allowing ScyllaDB to prioritize critical I/O operations. - NUMA Optimization: Minimizing cross-NUMA traffic ensures consistent memory access times, preventing latency spikes from remote memory fetches.
- Backpressure and Load Balancing: ScyllaDB incorporates sophisticated internal mechanisms for backpressure and load balancing across cores and nodes, preventing any single component from becoming a bottleneck and ensuring graceful degradation under overload.
Production Gotchas & Troubleshooting
-
Misconfigured CPU Pinning/Isolation:
- Symptom: Unpredictable latency, lower than expected throughput, high CPU steal time, or
topshowing ScyllaDB processes jumping between cores. - Cause: ScyllaDB relies on dedicated CPU cores. If the OS scheduler is allowed to move ScyllaDB threads or if other processes contend for the same cores, performance degrades.
- Fix: Ensure
isolcpusorcpusetkernel parameters are correctly configured in/etc/default/grub(or equivalent) to isolate cores for ScyllaDB. Usetuned-adm profile scyllaon RHEL/CentOS or manually configureirqbalanceto avoid routing interrupts to ScyllaDB cores. Verify withlscpu -eandtaskset -cp <pid>.
- Symptom: Unpredictable latency, lower than expected throughput, high CPU steal time, or
-
NUMA Misalignment:
- Symptom: High
numa_hitandnuma_missmetrics innumastat -m, orperfshowing significant remote memory accesses. - Cause: ScyllaDB processes are not correctly bound to their respective NUMA nodes, or memory is being allocated from a remote node.
- Fix: ScyllaDB typically handles NUMA binding automatically. Ensure
numactlis installed. Check ScyllaDB logs for NUMA-related warnings. Verifyscylla.yamlsmpandmemorysettings are appropriate for your hardware. For example, if you have two NUMA nodes,smpshould be half your total cores, and memory should be half your total RAM.
- Symptom: High
-
Insufficient
io_uring/ AIO Resources:- Symptom: I/O queue depth warnings in logs, high disk I/O latency despite fast storage,
iostatshowing highawaittimes. - Cause: The kernel's
io_uringor AIO limits are too low, or the underlying storage is saturated. - Fix: For
io_uring, ensure kernel is recent enough. For older AIO, increasefs.aio-max-nrin/etc/sysctl.conf(e.g.,fs.aio-max-nr = 1048576). Monitor disk I/O metrics closely. Consider faster storage (NVMe) or more I/O paths.
- Symptom: I/O queue depth warnings in logs, high disk I/O latency despite fast storage,
-
Network Stack Contention (without DPDK/XDP):
- Symptom: High network latency, packet drops,
netstat -sshowing errors, especially on high-throughput nodes. - Cause: The kernel's default network stack can become a bottleneck under extreme packet rates, leading to context switches and buffer bloat.
- Fix: For critical, high-performance deployments, consider enabling DPDK or XDP for ScyllaDB. This requires specific NICs and kernel modules. Otherwise, ensure network buffer sizes are tuned (
net.core.rmem_max,net.core.wmem_max,net.ipv4.tcp_rmem,net.ipv4.tcp_wmem).
- Symptom: High network latency, packet drops,
-
Over-provisioning/Under-provisioning Cores:
- Symptom: Low CPU utilization on some cores, high load on others, or overall poor performance.
- Cause:
scylla.yamlsmpsetting does not match the available isolated cores, or the workload is unevenly distributed. - Fix: Set
smpto the number of isolated cores dedicated to ScyllaDB. Leave at least one core for the OS and background tasks. Ensure your data model and partition keys lead to an even distribution of data and requests across nodes and cores.
Frequently Asked Questions
-
Why does ScyllaDB use C++ instead of a garbage-collected language like Java or Go? ScyllaDB uses C++ to achieve deterministic performance and fine-grained control over system resources. This avoids the unpredictable latency spikes caused by garbage collection pauses (common in Java/Go) and allows for direct memory management, NUMA awareness, and user-space I/O, which are critical for sub-millisecond P99 latency at scale.
-
How does ScyllaDB handle concurrency without traditional locks or OS threads? ScyllaDB employs a "shared-nothing, thread-per-core" model. Each CPU core runs a dedicated Seastar thread with its own local memory, network, and I/O queues. Concurrency within a core is managed via cooperative multitasking using futures and continuations. Inter-core communication occurs via explicit message passing, avoiding shared memory contention and locks.
-
What is
io_uringand why is it important for ScyllaDB's performance?io_uringis a modern Linux kernel interface for asynchronous I/O. It allows ScyllaDB to submit and complete I/O requests directly from user-space without system calls or context switches for each operation. This significantly reduces I/O overhead, improves throughput, and lowers latency compared to traditional kernel-buffered I/O or older AIO interfaces. -
How does ScyllaDB ensure data locality and NUMA awareness? ScyllaDB's Seastar engine is designed to be NUMA-aware. It partitions its memory pools and I/O queues such that each Seastar thread (pinned to a CPU core) allocates and accesses memory primarily from its local NUMA node. This minimizes slower cross-NUMA node memory access, which is crucial for consistent high performance on multi-socket servers.
-
Can I run other applications on the same server as ScyllaDB? While technically possible, it is strongly discouraged for production ScyllaDB deployments. ScyllaDB is designed to consume nearly all available CPU, memory, and I/O resources on its dedicated cores for optimal performance. Running other applications on the same server, especially on the same isolated cores, will lead to resource contention, unpredictable latency, and degraded ScyllaDB performance. Dedicated hardware or isolated virtual machines with CPU pinning are recommended.
Free In-Browser Developer Tools
Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.
Related Articles

ScyllaDB vs Apache Cassandra in 2026: P99 Latency, C++ Seastar & TCO Benchmarks
Comprehensive guide covering scylladb vs apache cassandra in 2026: p99 latency, c++ seastar & tco benchmarks with production-grade architecture and code examples.
Read more
Kafka vs Redpanda in 2026: Thread-per-Core Architecture, Zero-Disk Cache & P99 Latency Benchmarks
Comprehensive guide covering kafka vs redpanda in 2026: thread-per-core architecture, zero-disk cache & p99 latency benchmarks with production-grade architecture and code examples.
Read more
PostgreSQL 17 Query Optimization: Execution Plans, Memory Tuning & EXPLAIN ANALYZE
Comprehensive guide covering postgresql 17 query optimization: execution plans, memory tuning & explain analyze with production-grade architecture and code examples.
Read more