What a Cassandra benchmark is for
A Cassandra benchmark answers a capacity question: how many operations a cluster can serve, at what latency, for a specific workload, on specific hardware. Apache Cassandra is designed for high write throughput and horizontal scaling across commodity nodes, so most benchmarking questions concern the other side of that design: read latency under load, the number of nodes needed to hold a latency target, and the cost of that node count.
Because Cassandra’s behaviour depends heavily on the workload, a benchmark result is only meaningful together with its configuration. The same cluster can look ten times faster or slower depending on the read/write mix, the size of the working set relative to memory, the concurrency of the load generator and the consistency level. The sections below cover what to measure and how to interpret it. The final section summarises the benchmarks rENIAC has published for its Data Engine, with the configuration behind each number.
What a Cassandra benchmark should measure
- Throughput — operations per second sustained over the run, split by operation type (reads, writes and, where relevant, scans or batches).
- Latency distribution — not the average alone but p50, p95, p99 and, for latency-sensitive services, p99.9, measured at the client.
- Workload definition — read/write ratio, row and payload size, number of columns, key distribution and access skew.
- Concurrency — client threads or in-flight requests, number of client machines, and whether the load is closed-loop or rate-limited.
- Dataset — rows populated, bytes per node, and how the working set compares with DRAM and the operating-system page cache.
- Cluster topology — node count, replication factor, consistency level, driver and native-protocol version, network bandwidth.
- Resource use — CPU utilisation, JVM garbage-collection pauses, disk IOPS and queue depth, network throughput on every node and on the clients.
- Background activity — compaction, repair and hinted handoff running during the measurement window.
- Scaling — how throughput and tail latency change as concurrency and node count change.
Throughput versus latency
Throughput and latency are linked by the concurrency of the load. For a closed-loop load generator with N threads, each waiting for a response before issuing the next request, throughput is approximately N divided by the mean latency. Adding threads raises throughput until a resource saturates, usually CPU on the coordinator nodes or disk on the replicas. Beyond that point throughput stays flat while latency grows, because additional requests are simply queueing.
A useful benchmark therefore reports throughput at a stated latency, not either number alone. “50,000 reads per second” without a latency is a saturation figure; “50,000 reads per second with p99 below 10 ms” is a capacity figure that can be compared with a service-level objective. Plotting latency against offered throughput for several thread counts shows the knee of the curve, which is the operating point most production clusters should stay below.
Average latency versus tail latency
Cassandra’s average latency is usually low. The problems are in the tail. Garbage-collection pauses on the JVM, compaction competing for disk and CPU, read repair, reads that must merge several SSTables, and slow replicas that the coordinator waits for at higher consistency levels all produce occasional slow requests. They barely move the mean but define p99 and p99.9.
Service-level agreements are written at percentiles, so the benchmark must report them. p95 shows what one request in twenty experiences; p99 shows one in a hundred. A page that issues 50 queries in parallel has roughly a 40 percent chance of including at least one p99 response, so the database tail becomes the page’s typical behaviour. rENIAC’s 2017 white paper gives the customer-side view: a media company met its 7–8 ms target only at the 75th percentile, with 35 ms at the 95th and 60 ms at the 98th, and that gap was the actual problem to solve.
Measurement also matters. Closed-loop tools stop issuing requests while waiting on a slow response, so a pause in the server is recorded as one slow request instead of the many requests that would have arrived during the pause. This is coordinated omission; it makes tails look better than they are. Where possible, run with a fixed target rate as well as with a thread count, and check that the two agree.
Read/write workload mix
Cassandra writes are cheap: a write appends to the commit log and a memtable, and the expensive work of compaction happens later. Reads can be expensive: a read may consult the memtable, the bloom filters and key caches of several SSTables, and then disk, before the coordinator reconciles replica responses. Write-heavy benchmarks therefore flatter Cassandra and read-heavy ones expose its limits. Real applications are commonly read-dominated: rENIAC’s 2020 VMware article cites read/write ratios between 5:1 and 500:1 for large data-centric applications, and its AWS test schema for an inventory-management system uses 96 percent reads.
The mix changes results by large factors. In rENIAC’s virtualised 2020 test, the same Data Engine deployment gave 20x the baseline throughput on a 100 percent read workload, 7.4x at 90 percent reads and 3.4x at 80 percent reads. A benchmark that does not state its mix cannot be compared with anything.
Write-heavy runs must be long enough to include compaction. The same article attributes almost 65 percent of CPU cycles in a write-heavy workload to compaction and compression; a ten-minute run that finishes before the first major compaction measures a cluster that does not exist in production.
Concurrency
Concurrency has two parts: the number of client threads or in-flight requests, and the number of client machines. Either can bottleneck a benchmark before the database does. A single load generator saturating its own CPU or network interface produces a plateau that looks like a database limit. rENIAC’s tests used dedicated client servers, and the 2020 virtualised setup drove three database nodes from two client virtual machines so that the client was not the limiting resource.
Sweep concurrency rather than fixing it. Results at 16, 64, 256 and 1,024 threads describe a cluster; a single point does not. Inside Cassandra, watch the pending and blocked counts in the thread-pool statistics: a growing read-stage queue means the node is already past its knee.
Dataset size
The relationship between the dataset and memory decides which subsystem is being measured. A dataset that fits in the page cache benchmarks CPU and network. A dataset several times larger than memory benchmarks disk and the read path. Both are legitimate, but they are different tests, and confusing them is the most common reason two teams get incompatible numbers from the same hardware.
Report the population size and the size of the hot set, together with the key distribution. Uniform access spreads load; Zipfian or otherwise skewed access concentrates it on a few partitions, which is closer to most production traffic and is where caching pays off. rENIAC’s Amazon Keyspaces material describes exactly this case: skewed access causing throttling on frequently accessed keys.
Data per node is the other side of dataset size. A common rule of thumb, discussed in rENIAC’s data-density article, limits each Cassandra node to around 2 TB because read latency degrades as data per node grows. rENIAC’s 2020 data-density test compared a bare-metal baseline holding 574 GB with a Data Engine deployment holding 5.9 TB under an 80:20 mix at 20 KB per request, and reported 10.3x more data with 2.8x higher throughput. The point of that test is that capacity per node is a benchmark outcome, not only a configuration input.
Node count and scaling behaviour
Cassandra’s write path scales close to linearly with node count. Reads scale less predictably, because a read at consistency level QUORUM involves a majority of replicas and the coordinator’s work grows with the replication factor. A benchmark on a single node with replication factor 1 measures the storage engine; it says little about a nine-node cluster at replication factor 3, where the same request touches three machines and a network hop.
State the node count, the replication factor and the consistency level for every run, and keep them constant when comparing configurations. rENIAC’s 2018 lab test used one Cassandra node and one Data Proxy node to isolate the proxy’s effect; its 2020 vSphere test used a three-node Cassandra cluster and a three-node Data Engine cluster. The two results are not directly comparable with each other, and neither claims to be.
To characterise scaling, run the same workload against 3, 6 and 9 nodes and plot throughput at a fixed latency target. Linear scaling means the cost per operation stays constant as the cluster grows; sub-linear scaling is the signal that the read path, not the hardware count, is the constraint.
CPU and resource utilisation
Throughput and latency describe the service; utilisation explains it. Collect CPU utilisation, garbage-collection time, disk IOPS and queue depth, and network bytes on every node and every client throughout the run. A cluster delivering its target throughput at 90 percent CPU has no headroom for compaction, repair or a traffic spike; the same throughput at 40 percent CPU is a different result.
In rENIAC’s analysis of the Cassandra read path, storage and network I/O alone can consume about 60 percent of CPU in a read-heavy workload, and that I/O work is what its Data Engine offloads. A fair comparison of an accelerated and an unaccelerated cluster therefore reports utilisation on the database nodes as well as client-side throughput: the goal is fewer nodes at the same service level, and node count follows utilisation.
Why benchmark methodology matters
- Change one variable at a time. Compare configurations on identical hardware, software versions, schema, data and client settings. rENIAC’s published comparisons run the same cassandra-stress workload with and without the Data Engine in the path and change nothing else.
- Populate first, then measure. Reads against an empty or half-populated table hit bloom filters and return quickly. Write the full dataset, let compaction settle, then run the read phase.
- Warm up, then run long enough. Discard the first minutes; JIT compilation and cache warm-up distort them. Run long enough to include garbage-collection and compaction cycles.
- Repeat. Report the median of several runs and the spread between them. A single run cannot distinguish a real difference from noise.
- Measure at the client. Server-side latency excludes network time and queueing in the driver. The application experiences client-side latency.
- Use a standard tool and publish its parameters. cassandra-stress ships with Apache Cassandra and is the tool used in rENIAC’s own tests and in the Data Engine documentation; YCSB is the common alternative. The exact command line, thread count, distribution and row size are part of the result.
- Document the environment. Instance types, CPU model, memory, storage type, network bandwidth, Cassandra version, native-protocol version, replication factor, consistency level, compaction strategy and compression settings. Without them the result cannot be reproduced or compared.
Common mistakes when interpreting benchmark results
- Comparing averages. Two systems with the same mean latency can differ by an order of magnitude at p99, which is where the SLA lives.
- Reading a saturation number as capacity. Peak operations per second with no latency attached says where the cluster falls over, not where it can run.
- Ignoring the workload mix. A 20x result at 100 percent reads and a 3.4x result at 80 percent reads describe the same system; quoting one without the mix misrepresents it.
- Benchmarking a cache-hot dataset and generalising. If the data fits in memory, the disk was never tested.
- Letting the client set the ceiling. A saturated load generator produces a flat line that is easy to misattribute to the database.
- Comparing different consistency levels or replication factors. ONE against QUORUM, or replication factor 1 against 3, is not a like-for-like comparison.
- Stopping before background work starts. Short runs skip compaction, repair and garbage collection, the main sources of tail latency in production.
- Forgetting cost. The practical question is usually service level per dollar. Halving latency by doubling node count is a different outcome from halving latency and node count together.
rENIAC Data Engine and Cassandra performance
rENIAC Data Engine (rDE) is a transparent proxy and cache for Apache Cassandra, DataStax Enterprise and Amazon Keyspaces. It sits between the application and the database nodes, speaks the Cassandra Query Language, serves reads from its own flash-backed store and forwards writes to the cluster while keeping its cache consistent. Applications are repointed at the Data Engine cluster address; no code, schema or driver changes are required. Deployments use FPGA-backed hardware, either Intel FPGA cards in standard servers or AWS F1 instances, or the software edition for AWS i3 instances, which the Data Engine documentation describes as implemented in C++ and not subject to JVM garbage-collection or compaction pauses.
In benchmark terms, the Data Engine changes the read path. Reads that would have consulted memtables, SSTables and several replicas are answered by the proxy, so read throughput per node rises and tail latency narrows, while the Cassandra nodes keep the write pipeline and compaction. The published results below show the pattern: the largest gains appear on read-dominated mixes and in p99 latency, and they shrink as the write share grows.
Benchmarks published by rENIAC
| Published | Environment | Workload | Result as published |
|---|---|---|---|
| April 2018 — lab test, Turbocharging your Cassandra DB with rENIAC Data Proxy | Three Supermicro 2U servers (2 × Intel Xeon E5-2650 v1, 64 GB DDR3, 220 GB SSD, Intel 10 GbE) on a 10 GbE switch: one client, one Cassandra 3.x node (native protocol v4), one Data Proxy node | cassandra-stress, 1 M read-only operations over a 100 K-row population, 4 KB per row in 10 columns | Client-side p99 and p99.9 latency reduced by 10x or more (the post’s throughput, p95 and p99 charts are not available on this site) |
| June 2020 — VMware blog, reposted as Accelerating Virtualized & Distributed Cassandra databases with FPGAs | vSphere 7: three Cassandra 3.11 nodes, three Data Engine VMs (Intel Arria 10 GX FPGA, 16 cores, 64 GB RAM, 1 TB pass-through NVMe, 10 Gbps), two client VMs | cassandra-stress at 100, 90 and 80 percent reads | Throughput 20x, 7.4x and 3.4x the baseline; p95 latency 36x, 8.6x and 1.4x lower, respectively |
| August 2020 — lab test, Density is a DBA’s Best Friend | Bare-metal Cassandra baseline against Cassandra fronted by Data Engine | 80:20 read/write, 20 KB payload per request; 574 GB held by the baseline against 5.9 TB with Data Engine | 2.8x higher throughput while holding 10.3x more data |
| December 2020 — customer deployment, Accelerating Cassandra Performance: An eCommerce Case Study | Product catalogue on Apache Cassandra (DataStax Enterprise) on Google Cloud; Data Engine as a managed cache holding just over 1 TB on a single instance | Product-catalogue transactions (mix not published) | 2.9x throughput, 2.6x lower mean latency, 10.8x lower p99 latency, 48 percent lower cloud-instance cost |
| April 2017 — white paper, Data Store Acceleration-as-a-Service on Amazon FPGA Instances | Digital media company’s user-personalisation store on Cassandra Community 2.1.13; Distributed Data Engine deployed as a network service | Production personalisation queries | Queries per node from 2,905 to 20,000–26,000 (9x throughput); 5–8 ms at the 99th percentile, against 7–8 ms at the 75th, 35 ms at the 95th and 60 ms at the 98th before (10x lower latency) |
Each figure is the ratio between the same workload run with and without the Data Engine in the configuration described, as published by rENIAC in the linked article or document. The full e-commerce case-study document is not available on this site.
Reading these results
The 2020 series is the most instructive for methodology: one deployment, three workload mixes, results ranging from 1.4x to 36x. Any Cassandra benchmark claim, from any vendor, should be read the same way, with the mix, the dataset and the node count attached.