Teams usually move from a relational database to Apache Cassandra for one of three reasons: write volume that a single primary cannot absorb, a need to run active in several regions, or data growth that makes sharding a relational database by hand look worse than adopting a database that partitions by design. All three are good reasons. But a migration that copies tables one-to-one and rewrites SQL as CQL almost always fails, because Cassandra is not a relational database with a different syntax. It has no joins, no foreign keys, no general transactions, and it rewards a schema built around queries instead of entities.

This article treats the migration as a project: whether to migrate at all, how to derive a Cassandra schema from your access patterns with a worked example, how to replace each relational feature you rely on, and how to move live data with no downtime using a backfill, change data capture, shadow reads and a reversible cutover.

Advertisement

Decide what should move

Cassandra is excellent at high-volume writes, predictable single-partition reads, linear scaling by adding nodes, and multi-datacenter replication with no single primary. It is poor at ad hoc queries, joins, aggregations across the whole dataset, and workflows that need atomic updates across many rows. Before committing, list your ten most frequent queries and your transactional invariants. If most queries are lookups by a known key and the invariants fit inside one partition, migrate. If the workload is reporting, exploratory SQL or multi-entity transactions such as double-entry ledgers, keep the relational database for that part.

Many successful migrations are partial: the high-volume event, session or timeline data moves to Cassandra, while billing and administration stay relational. Check one more thing: multi-partition ACID transactions (Accord) arrive in Cassandra 6.0, which at the time of writing is still pre-release, with 5.0 the production line. Plan as if you only have single-partition lightweight transactions.

The mental model shift

Relational design starts from entities and normalizes them; queries are joined at read time. Cassandra design starts from queries and stores each answer pre-joined, in a table whose partition key is what the query filters on and whose clustering columns are how it sorts. One logical entity often lives in several tables, one per access pattern, and the application writes all of them. Storage is cheap and writes are fast; reads that hit one partition are what stays fast as the cluster grows. The data modeling basics page covers the rules; here they are applied to a real schema.

Advertisement

Worked example: an order system

Take a typical order system in PostgreSQL or MySQL:

customers(customer_id PK, email UNIQUE, name, created_at)
orders(order_id PK, customer_id FK, status, total_cents, created_at)
order_items(order_id FK, line_no, sku, qty, price_cents, PK(order_id, line_no))

Step one is the access-pattern inventory, taken from query logs, not from memory. Suppose it yields four queries: Q1, fetch an order with its items by order id; Q2, list a customer's orders, newest first, paged; Q3, find a customer by email at login; Q4, list today's orders by status for a fulfilment dashboard. Each becomes a table:

CREATE TABLE shop.orders_by_id (            -- Q1: order header + items in one partition
  order_id    uuid,
  line_no     int,
  customer_id uuid  STATIC,
  status      text  STATIC,
  total_cents bigint STATIC,
  created_at  timestamp STATIC,
  sku text, qty int, price_cents bigint,
  PRIMARY KEY ((order_id), line_no)
);

CREATE TABLE shop.orders_by_customer (      -- Q2: newest first
  customer_id uuid,
  created_at  timestamp,
  order_id    uuid,
  status text, total_cents bigint,
  PRIMARY KEY ((customer_id), created_at, order_id)
) WITH CLUSTERING ORDER BY (created_at DESC, order_id ASC);

CREATE TABLE shop.customers_by_email (      -- Q3: unique lookup
  email text PRIMARY KEY,
  customer_id uuid
);

CREATE TABLE shop.orders_by_status_day (    -- Q4: bucketed to bound partition size
  status text, day date, bucket int,
  created_at timestamp, order_id uuid, total_cents bigint,
  PRIMARY KEY ((status, day, bucket), created_at, order_id)
);

The join between orders and items disappears because items are clustering rows inside the order's partition, with header fields stored once as STATIC columns. Q4 shows the partition-sizing discipline: (status, day) alone would put every pending order of a busy day in one partition, so a bucket column (for example hash(order_id) % 8) spreads it, and the dashboard reads all eight buckets in parallel. Do the arithmetic: at 400,000 orders a day and about 100 bytes per row, one unbucketed partition is roughly 40 MB a day; eight buckets bring that to about 5 MB each, comfortably below the commonly cited ceiling of around 100 MB per partition. Because status is part of the key, a status change is a delete plus an insert, so give this dashboard table a TTL and stop reading old days.

Replacing relational features

Each relational feature you depend on needs a deliberate replacement:

Relational featureCassandra replacementWhat you give up
JoinsDenormalized query tables written togetherStorage, and the work of keeping copies in step
Foreign keysApplication checks; orphans tolerated or cleaned by a jobDatabase-enforced integrity
Auto-increment idsuuid or timeuuid generated by the clientDense, ordered integers
UNIQUE constraintLightweight transaction: INSERT ... IF NOT EXISTSLatency (Paxos round trips) on that write
Multi-row transactionLogged batch across denormalized tables; single-partition LWTIsolation; batches are atomic, not isolated
Secondary access pathsAnother query table, or an SAI index (5.0+) for additional filtersFree-form WHERE clauses
COUNT, SUM, GROUP BYPre-aggregated tables, counters, or SparkAd hoc analytics

Uniqueness is the subtle one. Registering a new customer must claim the email first, with INSERT INTO customers_by_email (email, customer_id) VALUES (?, ?) IF NOT EXISTS, and only write the other tables if [applied] is true. See lightweight transactions for the cost and the rule of never mixing LWT and plain writes on the same row. Writing one order to several tables uses a logged batch, which guarantees that all the writes eventually apply, but not that a reader never sees some before others; the batch article explains why large or multi-partition batches are expensive.

Migration architecture

A zero-downtime migration: backfill the past, stream the present, verify, then cut overApplicationwrites + readsSQLRelational DBsource of truthchange logCDC pipelineDebezium -> Kafkaupserts, source tsbulk exportBackfill jobSpark / DSBulkCassandraquery tablesshadow readsComparatordiff, metricsread same keyEvery write carries the source change time as USING TIMESTAMP, so a slow backfill never overwrites a newer CDC write.Cut over reads first, then writes; keep the old database warm until rollback is no longer needed.
Backfill, CDC, shadow reads and cutover. The relational database stays the source of truth until the final step.

A safe migration runs four streams of work in a fixed order. First, start capturing changes from the source, so nothing written during the backfill is lost. Second, backfill historical rows. Third, run shadow reads that compare answers from both stores while the relational side still serves users. Fourth, cut over reads, then writes, keeping a rollback path open until the new system has run cleanly through a full business cycle.

Capturing changes before the backfill matters because the backfill takes hours or days. If you snapshot first and start CDC afterwards, every write in between is lost. Debezium reads the MySQL binlog or PostgreSQL logical replication slot and publishes row changes to Kafka with the source position and commit time; a consumer you write (or a Kafka Connect sink for Cassandra) turns each change into upserts against every query table that row feeds.

Backfill without losing writes

The backfill reads the source in primary-key ranges and writes the target tables. For simple one-to-one tables, DataStax Bulk Loader (dsbulk) loads exported CSV or JSON; for anything involving joins or reshaping, the Spark Cassandra Connector is the usual tool, because Spark can read the source over JDBC, join, and write each query table in parallel.

The critical detail is ordering. A backfill row and a CDC change for the same key can arrive in either order. Cassandra resolves conflicts by write timestamp (last write wins), so set the timestamp explicitly from the source data instead of letting the coordinator stamp the time of the copy:

# Backfill writer (Python driver). updated_at comes from the source row;
# the CDC consumer uses the change's commit time the same way.
stmt = session.prepare(
    "INSERT INTO shop.orders_by_customer "
    "(customer_id, created_at, order_id, status, total_cents) "
    "VALUES (?, ?, ?, ?, ?) USING TIMESTAMP ?")

for row in source_rows:                      # read in primary-key ranges
    ts_micros = int(row.updated_at.timestamp() * 1_000_000)
    session.execute_async(stmt, (row.customer_id, row.created_at, row.order_id,
                                 row.status, row.total_cents, ts_micros))

Now a stale backfill row carries an older timestamp than the CDC update that already landed, so it loses, and both writers become idempotent: replaying either one is harmless. Deletes need the same care, because a delete written with an older timestamp than an earlier backfill insert does nothing. Keep concurrency bounded (a semaphore around execute_async), throttle the source reads so production does not suffer, and record the last completed key range so a crashed backfill resumes rather than restarts.

Dual writes versus change data capture

The alternative to CDC is dual writes: the application writes to both databases. It needs no new infrastructure, but it is hard to make correct. If the second write fails, the stores diverge, and nothing replays it. If two requests race, they can apply in opposite orders in the two stores. Dual writes are tolerable when you also run a reconciliation job and use source timestamps as above; CDC is preferable because the database log is a single, ordered, replayable record of every change, including those made by batch jobs and manual fixes that never pass through your application.

Verification and shadow reads

Verification proves the copy is right before anyone depends on it. Use three layers. Counts: rows per key range in both stores, which catches gross loss cheaply. Sampled diffs: pick random primary keys from the source, read the corresponding rows from every query table, and compare field by field. Shadow reads: in production, serve each request from the relational database as usual, asynchronously issue the equivalent Cassandra query, and record mismatches as a metric with a sample of the differing keys. Run shadow reads for days, through peaks and batch jobs, and investigate every mismatch class until the rate is zero or fully explained (for example, CDC lag of a second or two).

Shadow reads also give you real latency numbers at real concurrency, under the read consistency level you plan to use. LOCAL_QUORUM for both reads and writes is the usual choice when the application expects to read its own writes within one datacenter.

Cutover and rollback

Cut over reads first, behind a feature flag, ramping from 1% of traffic to 100% while watching error rates and latency. Writes still go to the relational database, and CDC keeps Cassandra current, so rollback is flipping the flag back. Cutting over writes is the one-way door: once the application writes directly to Cassandra, the relational database stops being current. Keep it as a fallback by running reverse replication (a job that writes Cassandra changes back) for a few weeks, or accept that rollback after this point means replaying writes. Only when the new path has survived a month-end, a deploy and a node failure should you decommission the old database.

Failure modes

  • Copying the relational schema. A table per entity plus ALLOW FILTERING or secondary indexes to answer queries turns every read into a cluster-wide scan. Fix: one table per query.
  • Unbounded partitions. Keys like customer_id for a customer with millions of events, or status alone, grow forever and become slow to read, compact and repair. Fix: bucket by time or hash.
  • Deletes as a workflow. Queue-like tables that insert and delete constantly accumulate tombstones that slow reads until compaction purges them. Fix: TTLs and time-bucketed partitions you stop reading.
  • Lost writes during backfill. CDC started after the snapshot, or copies written with the coordinator's timestamp. Fix: CDC first, source timestamps always.
  • Read-modify-write races. Code that read, changed and saved a row inside a SQL transaction now silently loses concurrent updates. Fix: redesign as idempotent upserts, or use a conditional update for the few true compare-and-set cases.

Trade-offs

The migration buys horizontal scale, multi-region writes and stable latency at the cost of query flexibility, enforced integrity and transactions. Denormalization moves complexity from the database into application code: every new query may need a new table and a backfill. That is worth it for a few high-volume access patterns and rarely worth it for everything, which is why polyglot designs that keep a relational store for the flexible part are so common.

What to do next

  1. Export a week of query logs and rank queries by frequency; write down the ten that matter.
  2. Design one Cassandra table per query, with partition-size arithmetic for the largest key.
  3. List every UNIQUE, foreign key and transaction in the old schema, and choose its replacement from the table above.
  4. Stand up CDC from the source into Kafka before writing any backfill code.
  5. Build the backfill with USING TIMESTAMP from source data, bounded concurrency and resumable key ranges.
  6. Add shadow reads with a mismatch metric, and run them until mismatches are zero or explained.
  7. Cut over reads behind a flag, then writes, and keep the old database for rollback through one full business cycle.
Key takeaway: Do not translate tables; translate queries. Inventory real access patterns, build one Cassandra table per query with bounded partitions, and replace each relational guarantee deliberately: lightweight transactions for uniqueness, logged batches for multi-table writes, client-generated ids for sequences. Move data by starting CDC first, backfilling with source write timestamps so last-write-wins resolves conflicts correctly, verifying with counts, sampled diffs and shadow reads, then cutting over reads before writes while keeping a rollback path.