ScyllaDB and Apache Cassandra look identical from the outside. They speak the same CQL protocol, accept the same table definitions, use the same consistency levels and partition data with the same kind of hash ring. That sameness is why the comparison is hard: every difference that matters lives inside the node, in how it uses cores, memory and disks, and in how the cluster moves data as it grows.
This page explains both engines from first principles, shows where the architectural choices turn into latency, cost and operational differences, walks through a fair evaluation and a migration, and ends with a checklist. It assumes you already know you want a wide-column, partitioned database; if that is still open, read when to pick Cassandra first.
What they share
Both are log-structured merge-tree databases. A write goes to a commit log for durability and to an in-memory table, the memtable. When the memtable fills it is flushed as an immutable sorted file, an SSTable, and background compaction merges SSTables so reads touch fewer files and deleted data is eventually purged. If that cycle is unfamiliar, the LSM tree page covers it in detail.
Both are leaderless and masterless for data. A partition key is hashed to a token, the token decides which nodes hold replicas, and the client chooses a consistency level per request: ONE, QUORUM, LOCAL_QUORUM and so on. Replicas converge through hinted handoff, read repair and anti-entropy repair. Both offer lightweight transactions for compare-and-set, both support secondary indexes and materialized views with caveats, and both use gossip to share liveness.
So the data model advice is the same for both: one table per query, bounded partitions, controlled tombstones. The differences below change the constants, not the rules.
The real difference: how a node uses its hardware
Cassandra is written in Java. A node runs one JVM with a staged design: requests arrive on network threads and are handed to pools of worker threads for reads, mutations, compaction and so on. All threads share one heap, one set of caches and the operating system's page cache. This portable design has two known costs: threads contend for shared structures, and the garbage collector must periodically walk the heap. Modern collectors on JDK 17 keep pauses short, but tail latency still depends on heap sizing and allocation rate.
ScyllaDB is a C++ reimplementation built on Seastar, a framework that pins one thread to each core and gives it its own memory. Each of these shards owns a slice of the node's token range, its own memtables, its own cache and its own SSTables. Shards do not share data structures and do not lock; when one needs work done by another, it sends a message. Each shard runs a cooperative scheduler, with separate CPU and I/O shares for user queries, compaction, streaming and repair, so background work is throttled by priority rather than by a fixed megabytes-per-second setting. Disk access uses direct, asynchronous I/O and ScyllaDB's own row cache instead of the kernel page cache.
As a result ScyllaDB tends to get more throughput per node and flatter tail latency from the same hardware, especially with many cores and fast NVMe. The cost is that a shard is a hard boundary: a hot partition saturates one core while others idle, and a long synchronous task inside a shard, a reactor stall, delays everything queued on that core.
Driver routing: token-aware and shard-aware
With a prepared statement the driver knows which bound value is the partition key, hashes it and sends the request to a replica directly. That is token awareness and both engines benefit. ScyllaDB adds one more level: shard-aware drivers keep a connection to each shard of each node and send a request to the core that owns the token, which saves a cross-core hop. ScyllaDB also listens on a shard-aware port, 19042 by default next to the usual 9042, so a driver can choose which shard its connection lands on.
Drivers that are not shard-aware still work against ScyllaDB; they just give some of the advantage back. ScyllaDB maintains forks of the common drivers that add shard awareness, and the Python one installs as scylla-driver while keeping the cassandra import, so code is portable:
# pip install scylla-driver (a fork of the DataStax Python driver; same "cassandra" import)
from cassandra.cluster import Cluster, ExecutionProfile, EXEC_PROFILE_DEFAULT
from cassandra.policies import DCAwareRoundRobinPolicy, TokenAwarePolicy
from cassandra import ConsistencyLevel
profile = ExecutionProfile(
load_balancing_policy=TokenAwarePolicy(DCAwareRoundRobinPolicy(local_dc="dc1")),
consistency_level=ConsistencyLevel.LOCAL_QUORUM,
request_timeout=2.0,
)
cluster = Cluster(["10.0.0.11", "10.0.0.12"], execution_profiles={EXEC_PROFILE_DEFAULT: profile})
session = cluster.connect("telemetry")
# Prepared statements carry the partition key position, so the driver can hash it
# and pick the replica node -- and, on ScyllaDB, the exact shard (core) inside it.
insert = session.prepare(
"INSERT INTO readings (device_id, day, ts, value) VALUES (?, ?, ?, ?)"
)
session.execute(insert, ("dev-42", "2026-10-02", 1759390000, 21.7))Always prepare statements so routing information exists; multi-partition batches defeat routing because the coordinator must fan out.
Distributing data: vnodes against tablets
Cassandra assigns each node a number of random tokens on the ring, virtual nodes, set by num_tokens (16 by default since 4.0). Every table in a keyspace is spread across the ring the same way. Adding a node means it takes over token ranges and streams the matching data from existing replicas, one node at a time, and the cluster is only balanced if every node has similar capacity and the token allocation was planned well.
ScyllaDB kept vnodes for compatibility and then introduced tablets, enabled by default for new keyspaces from ScyllaDB 2025.1. A table is split into tablets, small ranges that can be moved independently. A load balancer splits tablets as tables grow, merges them as they shrink and migrates them between nodes and shards in the background, so a new node starts serving within the time it takes to move some tablets rather than after a full ring rebalance. Topology and schema changes go through Raft, which makes concurrent node additions safe. Tablets also allow clusters that mix instance sizes.
-- ScyllaDB 2025.1 and later: new keyspaces use tablets unless you opt out.
CREATE KEYSPACE telemetry
WITH replication = {'class': 'NetworkTopologyStrategy', 'replication_factor': 3};
-- Opt out (for example, if a feature you need is not supported with tablets in your
-- version). This cannot be changed later with ALTER; you would have to recreate it.
CREATE KEYSPACE legacy_counters
WITH replication = {'class': 'NetworkTopologyStrategy', 'replication_factor': 3}
AND tablets = {'enabled': false};
-- The table itself is ordinary CQL and works unchanged on Cassandra.
CREATE TABLE telemetry.readings (
device_id text, day text, ts bigint, value double,
PRIMARY KEY ((device_id, day), ts)
) WITH CLUSTERING ORDER BY (ts DESC);The catch is feature coverage. Earlier ScyllaDB documentation listed features that did not work in tablets keyspaces, including counters, change data capture and lightweight transactions, and that list has been shrinking release by release. Check the tablets page for your exact version before creating keyspaces, because the choice is fixed at creation: you cannot ALTER a keyspace into or out of tablets.
Compaction, storage formats and the write path
Compaction is where both engines spend most of their background I/O and where most operational pain comes from. Cassandra offers size-tiered, leveled and time-window compaction, and Cassandra 5.0 added the Unified Compaction Strategy, which can be tuned to behave like either tiered or leveled compaction and adapts as data grows. Cassandra 5.0 also added trie-based memtables and the BTI SSTable format, storage-attached indexes and vector search. UCS and the trie formats are opt-in, so an upgraded cluster keeps its old strategies until you change them. Release coverage credits UCS with supporting much denser nodes; treat that as a goal to test, not a guarantee.
ScyllaDB supports the same strategy names and adds incremental compaction, which compacts in fixed-size fragments so a large compaction does not temporarily need as much free space as the data it rewrites. Its compaction is driven by a backlog controller: the schedulers give compaction more CPU and I/O share when the backlog grows and less when queries need it, instead of you guessing a throughput cap. The write amplification page explains why the strategy choice matters for SSD wear and space headroom; the trade-offs apply to both engines.
A worked evaluation and migration
Suppose a team runs a 12-node Cassandra 4.1 cluster for device telemetry: partitions keyed by device and day, about 60,000 writes per second and 20,000 reads per second at peak, a 99th percentile read target of 10 ms, and roughly 40 TB of data before replication. Licence and node costs are rising and they want to know whether ScyllaDB would let them run fewer nodes.
Step one is a fair benchmark. Use production-shaped data, the same schema and the same consistency level on both systems, on identical instance types. Load enough data that it does not fit in memory, then run a mixed workload at a fixed rate for hours so compaction reaches steady state. Comparing maximum throughput alone hides the interesting part; compare latency percentiles at the rate you actually need, and then at 1.5 times it.
# Same load generator, same schema, both clusters. Run each phase long enough
# for compaction to reach steady state (hours, not minutes).
cassandra-stress write n=200000000 cl=LOCAL_QUORUM \
-rate threads=256 -node 10.0.0.11,10.0.0.12,10.0.0.13 -log file=write.log
cassandra-stress mixed ratio\(write=1,read=3\) duration=2h cl=LOCAL_QUORUM \
-rate threads=256 fixed=60000/s -node 10.0.0.11 -log file=mixed.log
# Compare p99 and p999 at the SAME fixed rate, not maximum throughput alone.Step two is the capacity decision. Fewer nodes only save money if they also hold the data with compaction headroom and survive a node loss; fewer, larger nodes move more data per failure, so time a node replacement as part of the test.
Step three is migration without downtime: dual writes from the application, a historical backfill with SSTable streaming or ScyllaDB's Spark-based migrator, validation by row counts per token range and sampled full-row comparisons, then a read cut-over one service at a time. Keep dual writes until rollback is no longer needed.
Failure modes on each side
- Hot partitions. Both suffer, but on ScyllaDB one hot key pins one core at 100 percent while the node looks idle overall. Watch per-shard load, not node averages, and salt or bucket keys that attract heavy traffic.
- Large partitions and tombstones. Both engines suffer. Bucket by time, set TTLs deliberately and watch large-partition warnings.
- Garbage collection pauses. Cassandra-specific. An undersized heap or a very high allocation rate shows up as periodic latency spikes and, at worst, nodes marked down by peers. Measure pause times, not just heap usage.
- Reactor stalls. ScyllaDB-specific. A shard blocked by a long task delays all its queued requests. Stalls are logged with backtraces; recurring ones usually trace to a giant partition, a huge batch or an unusual query pattern.
- Compaction backlog and full disks. Common to both. A growing backlog means too little capacity for the write rate or the wrong strategy.
- Skipped repair. Common to both. Without regular repair, deleted data can reappear once tombstones expire. Schedule repairs inside the gc grace period on every table.
Licensing and ecosystem
| Aspect | Apache Cassandra | ScyllaDB |
|---|---|---|
| Licence | Apache 2.0, governed by the ASF | Source-available from 2025.1; free tier up to 50 vCPUs and 10 TB total storage per organization. OSS 6.2.x and earlier remain AGPL but get no new features |
| Commercial support | Several vendors and managed services | ScyllaDB the company, plus its managed cloud service |
| Extra APIs | CQL | CQL plus a DynamoDB-compatible API (Alternator) |
| Tooling | nodetool, cassandra-stress, broad third-party ecosystem | nodetool-compatible commands, ScyllaDB Manager and Monitoring Stack |
| Tuning surface | JVM heap and GC, thread pools, compaction throughput | Mostly automatic schedulers; fewer knobs, less to tune and less to override |
The licence change matters for some organizations more than any benchmark. If a permissive, foundation-governed licence is a requirement, Cassandra is the answer. If the cluster stays below the free tier, or you would pay for support anyway, it may not matter.
How to choose
| Situation | Leaning | Why |
|---|---|---|
| Large cluster where node count drives cost | ScyllaDB | Per-core design usually means fewer nodes for the same load; verify with your workload |
| Strict p99 or p999 targets | ScyllaDB | No GC and per-shard scheduling flatten tail latency |
| Open governance and licence are requirements | Cassandra | Apache licence; many vendors |
| Heavy use of features with gaps on tablets in your version | Cassandra, or ScyllaDB with vnode keyspaces | Check the tablets limitations first |
| Team already expert in Cassandra on the JVM | Either; measure | Operational skill transfers; the data model does not change |
| Wants the newest Cassandra 5.0 features such as SAI or vector search | Cassandra 5.0, or compare with ScyllaDB's equivalents | Feature sets diverge; check exact semantics before porting |
For sharding background that applies to both, see sharding strategies compared, and for compaction tuning on the Cassandra side see Cassandra compaction strategies.
What to do next
- Write down the workload in numbers: reads and writes per second at peak, partition size distribution, latency percentiles you must hold and total data size.
- Check the licence question with whoever owns procurement before any benchmark; it can decide the matter alone.
- If you would use ScyllaDB tablets, confirm your feature list against the tablets limitations for the exact release you would deploy.
- Run cassandra-stress or your own replay against both engines on identical hardware, at a fixed rate, for hours, with data larger than memory.
- Time a node replacement on each system at the node size you would actually run.
- Audit the schema for hot and large partitions; fix them first, because neither engine forgives them.
- If you migrate, plan dual writes, backfill, validation per token range and a reversible read cut-over.