Cassandra's counter type looks like the easiest feature in CQL: declare a column as counter and write SET views = views + 1. Underneath it is the one part of Cassandra's write path that is not blind. Every ordinary write is a timestamped cell that replicas can apply in any order, which is why writes are cheap and retries are safe. A counter increment cannot work that way, because two increments must both count rather than one overwriting the other. Cassandra solves this with a per-replica sharded value and a read-before-write on one replica, and that design explains every surprising property counters have: they are slower than normal writes, they cannot be retried safely, and they are not exact.

This article explains the design from the storage format up, using the Cassandra 4.1 source as the reference: one increment through the cluster, why a timed-out increment may have been applied, schema rules, deletes, hot keys, configuration, and what to use when counters are the wrong tool.

Advertisement

The rules the schema enforces

Counters have their own table discipline. A counter column cannot be part of the primary key, and every non-primary-key column in a table that holds counters must itself be a counter. You cannot INSERT into a counter or set it to an absolute value; the only write is an increment or decrement by a delta. Counters do not support TTLs or client-supplied timestamps, and they cannot be indexed.

CREATE TABLE page_views (
    page_id  text,
    day      date,
    views    counter,
    uniques  counter,
    PRIMARY KEY ((page_id), day)
);

UPDATE page_views SET views = views + 1, uniques = uniques + 1
 WHERE page_id = 'home' AND day = '2026-09-30';

-- Counter updates in a batch must use a counter batch, and cannot mix in regular writes
BEGIN COUNTER BATCH
  UPDATE page_views SET views = views + 3 WHERE page_id = 'docs' AND day = '2026-09-30';
  UPDATE page_views SET views = views + 1 WHERE page_id = 'blog' AND day = '2026-09-30';
APPLY BATCH;

If you need metadata next to a count, put it in a separate regular table with the same key.

The storage format: a counter is a set of shards

A counter cell does not store a single number. It stores a counter context: a list of shards, each a tuple of (counter id, logical clock, count). Every node has a counter id, a UUID it uses to identify the shard it owns. The value you read is the sum of the counts across all shards.

When a replica increments a counter, it only ever changes its own shard: it bumps the logical clock and adds the delta to its own count. When two versions of the same counter meet, during a read, in compaction, or through repair, shards with the same counter id are merged by keeping the one with the highest logical clock. Shards with different counter ids are kept side by side. Because each shard has exactly one writer and a monotonic clock, keep-the-highest is always correct, and merging is safe in any order and any number of times.

If this sounds like a CRDT, it nearly is. It resembles a grow-only counter per replica, with the twist that each shard's count can go up and down because the shard's owner is the only one updating it. The CRDT article covers the general theory; the important point here is that the merge is deterministic, so replicas that saw different histories converge once they exchange data.

The source also carries merge rules for older local and remote shard types, kept to read data written before the Cassandra 2.1 counter redesign. New writes create only global shards.

Advertisement

The write path, one increment at a time

ClientUPDATE ... c = c + 1Coordinatorpicks a leader replicaLeader replica (random live replica in local DC)1. take striped lock for the cell2. read local (clock, count): counter cache or disk3. clock = max(now_us, clock + 1); count += delta4. write global shard (own counter id)5. update cache, release lockforward deltaReplica Bstores A's shard as isReplica Cstores A's shard as isreplicate resultCounter cell on every replicashards: (A, clk, n) (B, clk, n) (C, clk, n)value = sum of counts; merge keeps highest clock per counter idack at CLOnly the leader does read-before-write. Other replicas receive a finished shard, which is why repair can reconcile it.
A counter increment. The coordinator forwards the delta to a leader replica, which does a locked read-before-write on its own shard and replicates the resulting shard, not the delta, to the other replicas.

The coordinator first picks a leader for the mutation: a random live replica in its local data centre, or the closest replica by snitch if the local data centre has none. If it picks itself it applies the increment locally; otherwise it forwards it. For DC-local consistency levels with no local replica alive, the write fails as unavailable. The coordinator checks up front that the requested consistency level can be met before forwarding anything.

On the leader, the increment runs as a small critical section. The leader takes striped locks covering the counter cells being changed. Cassandra 4.1 sizes the stripe set at concurrent_counter_writes times 1024, so unrelated counters rarely share a lock. If the lock cannot be taken within the counter write timeout, the write fails with a write timeout of type COUNTER. Holding the lock, the leader reads its own current (clock, count) for each cell, from the counter cache if present or from memtables and SSTables if not, computes the new clock as the larger of the current time in microseconds and the old clock plus one, adds the delta, and writes a new global shard under its own counter id.

Then it replicates. What goes to the other replicas is the resulting shard, the absolute (counter id, clock, count) tuple, not the delta. The other replicas do not read anything; they apply it like an ordinary write, and the merge rule makes it idempotent on their side. The coordinator acknowledges the client once enough replicas have acknowledged for the consistency level.

So every increment is a read on the leader, touching SSTables on a cache miss, and two concurrent increments on one leader must serialise on the lock, or both would read the same old count.

Worked example: three replicas, one shard each

Take RF=3 with replicas A, B and C, and a fresh counter. A client in the local data centre sends +5, and the coordinator picks A as leader. A reads nothing (blank), writes shard (A, t1, 5), and sends that shard to B and C. Everyone now holds {(A, t1, 5)} and the value is 5.

Next, a client sends +2 and the coordinator picks B. B's own shard is blank, so it writes (B, t2, 2) and replicates it. The cell is now {(A, t1, 5), (B, t2, 2)}, value 7. Then +3 goes to A again: A reads its own shard (A, t1, 5), writes (A, t3, 8), and replicates. The cell is {(A, t3, 8), (B, t2, 2)} and the value is 10.

Now suppose C was down during the last increment. C still holds {(A, t1, 5), (B, t2, 2)}, value 7. A read at QUORUM that includes C and A merges shard by shard: for counter id A it keeps t3 over t1, so the result is 10, and read repair or anti-entropy repair brings C up to date by the same rule. Nothing was double counted, because C never adds; it only takes the higher-clocked version of each shard.

Timeouts, retries and overcounting

The weak spot is the client. Suppose the leader applied the increment and replicated it, but the acknowledgement did not reach the client before the timeout, perhaps because a replica was slow at QUORUM. The client sees a write timeout and cannot know whether the increment happened. If it retries, and the first attempt had in fact been applied, the counter is now one too high. There is no idempotency key and no way to ask whether that increment landed: the new attempt is a new delta, applied again by whichever leader it reaches.

So counter increments are not idempotent, and a retry policy that treats them as idempotent will overcount under load. Drivers generally do not retry writes they cannot tell are idempotent; with the DataStax Python driver, a statement's is_idempotent flag defaults to false, and you should leave it that way for counters:

from cassandra.cluster import Cluster
from cassandra import WriteTimeout

session = Cluster(["10.0.0.1"]).connect("metrics")
incr = session.prepare(
    "UPDATE page_views SET views = views + ? WHERE page_id = ? AND day = ?")
incr.is_idempotent = False            # never let speculative execution or retries resend it

def record_view(page, day, n=1):
    try:
        session.execute(incr, (n, page, day))
    except WriteTimeout:
        # Outcome unknown. Choose explicitly: drop (risk undercount) or
        # retry (risk overcount). For view counts, dropping is usually right.
        metrics.counter_unknown.inc()

Pick the error you prefer per use case. Analytics counts usually accept a small undercount and drop on timeout. If you cannot accept either error, a counter is the wrong type, and the alternatives below apply. Also avoid speculative execution for counter statements for the same reason: a second in-flight copy is a retry by another name.

Consistency levels and reads

Counters use the same consistency levels as other writes, except ANY, which is not supported because a hint cannot perform the leader's read-before-write. The quorum reasoning from consistency levels applies to reads: writing and reading at LOCAL_QUORUM returns the most recent acknowledged value in that data centre. The leader's read of its own shard is local regardless of consistency level, which is safe because only that node ever writes its shard.

Lightweight transactions solve a different problem: linearizable compare-and-set at the cost of several Paxos round trips, the tool when a value must be exact, such as unique sequence numbers. Counters give cheap, commutative, approximately exact aggregation. An exact counter built from LWT on a high-traffic key will be throttled by contention long before a counter would suffer.

Deletes are one-way

You can delete a counter, but you should treat it as final. A delete writes a tombstone with a timestamp, while counter shards order themselves by their own logical clocks, and the two do not compose cleanly. Increments that race with the delete, or that arrive after it, can be lost or resurrect old shards once the tombstone is purged after gc_grace_seconds. The cassandra.yaml comment goes further: if you delete counters and rely on a low gc_grace_seconds, disable the counter cache, because a cached shard can outlive its tombstone.

The practical rule is never to reset a counter by deleting and incrementing it again. Put a time bucket in the key instead, as the page_views table above does with its day column, and let old buckets expire through deletion of whole partitions you no longer read.

Hot counters and performance

A single very hot counter, such as a global like count on a viral item, concentrates increments on one partition and, through the lock, on a serial section per leader. Throughput is bounded by how fast one replica can take a lock, read and write. The fix is the same as for any hot key: split the counter into N sub-counters and sum on read.

import random
SHARDS = 16

def incr_likes(item_id, n=1):
    s = random.randrange(SHARDS)       # spread load across 16 partitions
    session.execute(incr_likes_stmt, (n, item_id, s))

def read_likes(item_id):
    rows = session.execute(read_all_stmt, (item_id, list(range(SHARDS))))
    return sum(r.likes for r in rows)

Here the table key is ((item_id, shard)). Reads cost 16 partition lookups, so cache the sum for a few seconds on the read side.

The counter cache is the other lever. It keeps only the local (clock, count) tuple for each hot cell, not the whole counter, so it is cheap. With RF=1 a cache hit skips the read before write entirely; with RF above 1 it still shortens the time the lock is held. In 4.1 an empty counter_cache_size means automatic sizing, the smaller of 2.5 percent of heap and 50 MiB, and 0 disables it.

Configuration and metrics

Setting (4.1 names)Default in 4.1What it controls
concurrent_counter_writes32Threads for the counter mutation stage; counter writes read, so size like reads
counter_write_request_timeout5000msCoordinator wait for counter writes, including the leader's lock wait
counter_cache_sizeempty (auto)Memory for cached local shards; 0 disables
counter_cache_save_period7200sHow often cache keys are saved for warm restarts

Check these against your own cassandra.yaml: key names became duration strings in 4.1, and the development branch ships a different counter timeout. Watch counter write latency and timeouts separately from regular writes, pending counter mutation tasks, counter cache hit rate and client-side unknown-outcome counts. Rising timeouts on a stable workload usually mean a hot counter or slow disks under the leader's read; tracing one increment, which logs lock acquisition and cache hits, tells you which.

When not to use counters

Counters are right for high-volume approximate aggregates, such as page views, likes and usage meters, where a rare error of one is acceptable. They are wrong for money, inventory and anything needing exactly-once semantics or an audit trail. For those, write each event as an immutable row with a unique ID, so retries become harmless upserts, and aggregate asynchronously, or keep the balance in a transactional store. Read the write and read path article if you want to compare the blind-write path those designs use with the counter path above.

What to do next

  1. List every counter table and label each use as approximate (fine) or must-be-exact (move to an event log or a transactional store).
  2. Mark counter statements non-idempotent in the driver and disable speculative execution for them.
  3. Decide per use case whether a timeout is dropped or retried, and count unknown outcomes as a metric.
  4. Put a time bucket in counter keys instead of deleting and reusing counters.
  5. Split any counter taking thousands of increments per second into N sharded sub-counters and cache the sum.
  6. Confirm counter settings in your own cassandra.yaml and alert on counter write timeouts and counter stage backlog.
Key takeaway: A Cassandra counter is a set of per-replica shards, each written only by its owner and merged by keeping the highest logical clock, so replicas converge safely. The cost is a locked read-before-write on a leader replica for every increment, and the weakness is the client: a timed-out increment may have been applied, and retrying it overcounts. Use counters for approximate aggregates, bucket them by time instead of deleting, shard hot ones, and move anything that must be exact to an event log or a transactional store.