A counter column is the CQL type for a number you only ever change by a delta: page views, likes, API calls per customer per hour, bytes served per tenant. You declare it as counter, write SET n = n + 1, and Cassandra merges increments from many clients without you reading the old value first. It is the easiest way to keep a running total in Cassandra and also one of the easiest to misuse.

How a counter works inside the cluster, with per-replica shards, a leader replica that reads before it writes, and the reason a retried increment can count twice, is covered in Cassandra counters architecture, in depth. This article is about the layer above: how to design tables that use counter columns, how to cut the number of increments you send, how to deal with hot counters, how to give the business an exact number when counters are approximate, and how to operate and migrate counter tables. The examples use CQL and the DataStax Python driver against Cassandra 4.x and 5.0.

The rules, read as design constraints

Counter columns come with schema rules that shape every design decision, so it helps to state them as design consequences rather than as a list of errors.

RuleDesign consequence
A table with a counter column can contain only counters outside the primary keyNames, labels and other metadata live in a separate regular table with the same key.
A counter cannot be part of the primary keyYou can never query or sort by a counter's value; rankings need another structure.
The only write is an increment or decrement; no INSERT, no setting a valueYou cannot reset to zero or correct a value directly; plan resets into the key.
No TTL and no client-supplied timestampsOld counters do not expire; retention needs time in the key and explicit deletes or table rotation.
Counter updates batch only with other counter updatesCounter and regular writes for one event are two separate, non-atomic writes.
Deleting a counter is effectively permanentTreat a deleted key as dead; do not increment it again.
Increments are not idempotentA timed-out increment has an unknown outcome; exact totals need a second source of truth.

The last two rows are the ones teams forget. The documentation warns that reusing a counter after deleting it is unsupported, and the leader write path means a client cannot tell whether a timed-out increment landed. Every pattern below either works with those properties or adds something next to the counter to compensate.

Modelling: time in the key, one table per granularity

Event streampage views, API callsPre-aggregatorsum deltas per key, flush 5sCOUNTER BATCHCounter rollupsviews_by_minute, views_by_dayEvent log (regular table)event_id key, idempotentnightlyBatch recomputeSpark or a scan jobExact totalsregular bigintDrift monitorcounter minus exact, by keyMetadata tablepage_id, title, owner: same key, regular columnsFast path: counters for live display. Slow path: exact numbers for billing and reports.
A counter design that holds up in production: pre-aggregated increments into time-bucketed rollups for live numbers, an idempotent event log recomputed into exact totals, and a monitor comparing the two.

The usual mistake is one counter per thing, keyed by the thing: views counter keyed by page_id. It answers only all-time totals, it can never be reset, and it grows hotter as the thing becomes popular. Put time in the key instead, and keep one table per granularity you will read.

CREATE TABLE metrics.views_by_minute (
    page_id  text,
    day      date,
    minute   smallint,          -- 0..1439
    views    counter,
    bytes    counter,
    PRIMARY KEY ((page_id, day), minute)
);

CREATE TABLE metrics.views_by_day (
    page_id  text,
    month    text,              -- '2026-10'
    day      date,
    views    counter,
    bytes    counter,
    PRIMARY KEY ((page_id, month), day)
);

CREATE TABLE metrics.pages (    -- metadata cannot share a counter table
    page_id text PRIMARY KEY,
    title   text,
    owner   text
);

Partition sizing is now predictable. A minute partition holds at most 1,440 rows per page per day, and a day partition at most 31 rows per page per month. A dashboard reading today's minute series is a single-partition slice; a monthly report reads one partition per page. Retention becomes a matter of deleting whole old partitions, or of rotating to a new table per year and dropping the old one, which avoids leaving counter tombstones scattered across live partitions. The same bucketing reasoning is developed for regular tables in time-series modelling in Cassandra.

Both granularities should be written from the same flush, in one counter batch per partition where possible. A counter batch is not a logged batch, so it is not atomic across partitions; if a flush half-succeeds the minute and day tables disagree slightly, and the reconciliation path below is what catches that.

Send fewer increments: pre-aggregate and flush

Every counter increment costs a read on the leader replica and a lock on the cell, so the single most effective optimisation is to send fewer of them. If a page gets 2,000 views a second, sending 2,000 increments of one is wasteful when one increment of 2,000 every few seconds carries the same information. Aggregate in the application or the stream processor, then flush.

import threading, time
from collections import defaultdict
from datetime import datetime, timezone
from cassandra.cluster import Cluster
from cassandra.concurrent import execute_concurrent_with_args

session = Cluster(["10.0.0.1"]).connect("metrics")
incr_min = session.prepare(
    "UPDATE views_by_minute SET views = views + ?, bytes = bytes + ? "
    "WHERE page_id = ? AND day = ? AND minute = ?")
incr_min.is_idempotent = False          # never retried by the driver

pending = defaultdict(lambda: [0, 0])   # (page, day, minute) -> [views, bytes]
lock = threading.Lock()

def record(page_id, nbytes):
    now = datetime.now(timezone.utc)
    key = (page_id, now.date(), now.hour * 60 + now.minute)
    with lock:
        pending[key][0] += 1
        pending[key][1] += nbytes

def flush():
    global pending
    with lock:
        batch, pending = pending, defaultdict(lambda: [0, 0])
    params = [(v, b, page, day, minute) for (page, day, minute), (v, b) in batch.items()]
    results = execute_concurrent_with_args(
        session, incr_min, params, concurrency=50, raise_on_first_error=False)
    for (ok, result), p_ in zip(results, params):
        if not ok:
            log_unknown_outcome(p_, result)   # do not re-send: outcome unknown

def flusher(interval=5.0):
    while True:
        time.sleep(interval)
        flush()

Walk the numbers. Ten application instances each see 200 views a second on a popular page. Without aggregation that is 2,000 counter writes per second against one cell, each serialised on the leader's lock. With a five-second flush, it is ten increments every five seconds, two per second, for the same total. The cost is a loss window: if an instance dies, up to five seconds of its counts vanish, and on shutdown you must call flush() before exiting. For live dashboards that is almost always an acceptable trade; for billing it is not, which is why billing should never rely on counters alone.

Failed flushes are logged, not retried. A timed-out increment may already have been applied, so re-sending it risks double counting. The log of unknown outcomes becomes an input to reconciliation.

Hot counters: splitting inside and across partitions

Pre-aggregation solves most hot counters. When it cannot, for example when thousands of independent clients write directly, split the counter. There are two levels, and they fix different problems.

  • Split inside the partition. Add a shard tinyint clustering column and have each writer pick a random shard from 0 to 15. Sixteen cells means sixteen independent locks on the leader, so lock contention drops, and a read sums sixteen rows from one partition. But all sixteen rows still live on the same replicas, so this does not spread disk or CPU load across the cluster.
  • Split across partitions. Put the shard in the partition key, PRIMARY KEY ((page_id, day, shard), minute). Now the sixteen sub-counters land on different replica sets and the write load spreads, but a read becomes sixteen partition reads, best issued concurrently and summed in the client.

Choose the shard count from measured load, not guesswork: start with the write rate per key that your cluster handles comfortably in a load test, divide the peak rate by it, and round up to a power of two. The general technique, and its accuracy trade-offs, is covered in distributed counter architecture. The key point for Cassandra is that partition-level splitting is the one that addresses replica hotspots described in Cassandra partitioning.

Reading counters and building rankings

Counters are cheap to read by key and impossible to query by value. You cannot index a counter, and because a counter cannot be a clustering column, you cannot ask for the top ten pages by views. Leaderboards need a separate structure: a periodic job reads the day's counters, computes the ranking and writes it into a regular table clustered by score, such as PRIMARY KEY ((board, day), score, page_id) with descending order on score. Readers query that table; it is a few minutes stale, which is normal for rankings.

Reading at a weaker consistency level is usually fine for display. If the dashboard reads at ONE and one replica missed an increment, the number may be slightly low until repair or read repair catches it up, and the next refresh will usually be right. Reads at QUORUM cost more and still cannot recover an increment whose outcome was unknown.

Exact totals next to approximate counters

Counters are approximate in two ways: increments can be lost when a client gives up on an unknown outcome, and they can be doubled when a client retries one that had in fact succeeded. When a number feeds billing, quotas or published reports, keep an exact source of truth beside the counter.

CREATE TABLE metrics.view_events (
    page_id  text,
    day      date,
    event_id timeuuid,         -- generated once by the producer
    bytes    bigint,
    PRIMARY KEY ((page_id, day), event_id)
) WITH default_time_to_live = 7776000;   -- 90 days

CREATE TABLE metrics.views_exact (
    page_id text,
    day     date,
    views   bigint,
    bytes   bigint,
    PRIMARY KEY (page_id, day)
);

Writing an event with a producer-generated event_id is idempotent: a retry rewrites the same cell, so it can be retried freely. A nightly job, typically Spark with the Cassandra connector or a token-range scan, counts each partition of view_events and writes views_exact with a plain INSERT. A drift monitor then compares views_exact with the counter's day value and alerts if the gap exceeds, say, 0.1 percent for any key. Persistent drift points to retry policies that resend increments or to frequent timeouts.

If you need an exact value updated in real time rather than nightly, a counter is the wrong type; a regular bigint updated with a lightweight transaction gives compare-and-set at a much higher cost per write, as explained in lightweight transactions.

Operations: caches, resets and migration

Operating counter tables is mostly about the counter write path. Watch write latency and timeouts for the counter tables specifically, with nodetool tablestats metrics and your metrics exporter, because counter timeouts mean unknown outcomes rather than just slow requests. The counter cache, sized in cassandra.yaml, keeps recently used local shard values in memory so the leader's read-before-write avoids SSTables; check its hit rate with nodetool info before raising it. Run repair on counter tables like any other table, because repair is how a replica that missed shards catches up.

To reset or rebase, never try to decrement to zero under live traffic, since concurrent increments make the result unpredictable. Change the key instead: add an epoch or period column to the partition key and start writing a new epoch. To migrate counters to a new table or cluster, stop writers or switch them to the new table, read each final value from the old one and apply it as a single increment to a fresh key in the new one, then verify totals against the exact table. Copying counter SSTables between clusters is a specialist operation; the read-and-increment route is slower but easy to verify.

Trade-offs

ApproachWrite costExact?Use for
Counter columnRead on the leader plus a lockNoLive dashboards, rate displays, rough usage
Pre-aggregated counterFew incrementsNo, plus a loss windowHigh-volume metrics
Event log plus batch recomputePlain idempotent write per eventYes, after the batchBilling, reports, audits
LWT on a bigintPaxos round trips per writeYes, in real timeLow-rate quotas, balances

Most systems should combine the first and third rows: counters for speed and an event log for truth, with a monitor showing how far apart they are.

What to do next

To apply this to your own tables:

  1. List every counter table and write down which reader needs exact numbers; give those readers an event log and a recompute job.
  2. Put time in every counter's partition key and create one table per read granularity.
  3. Move metadata out of counter tables into regular tables with the same key.
  4. Add a pre-aggregation buffer in front of the busiest counters and flush on shutdown.
  5. Set is_idempotent = False on counter statements and log unknown outcomes instead of retrying.
  6. Load test the hottest key and choose a shard count from the measured safe rate per key.
  7. Build the drift monitor and alert when counters and exact totals diverge.
Key takeaway: Counter columns give cheap running totals but cannot be set, sorted, indexed, expired or safely retried. Design around that: put time in the partition key with one table per granularity, keep metadata in a separate table, pre-aggregate increments and flush every few seconds, split hot counters across partitions when one key is too busy, build rankings in a separate table, and keep an idempotent event log recomputed into exact totals with a drift monitor whenever the number feeds billing or reports.