Databases 35 min read

Spring Boot + Cassandra for High-Write Workloads: Data Modeling, Consistency Tuning & Pitfalls

This guide covers integrating Cassandra with Spring Boot for high-write scenarios, detailing data modeling with partition/clustering keys, consistency level selection (LOCAL_QUORUM), query best practices, performance tuning (compaction, compression, JVM), and production pitfalls like hot partitions and tombstone storms.

Xiaolin Talks Programming
Xiaolin Talks Programming
Xiaolin Talks Programming
Spring Boot + Cassandra for High-Write Workloads: Data Modeling, Consistency Tuning & Pitfalls

Why Cassandra for High-Write Scenarios

Device metrics, logs, and transaction flows are write-heavy and ever-growing. MySQL hits three walls: write amplification (clustered index, secondary indexes, redo/undo/binlog), painful resharding (dual-write, sync, verification), and cross-region replication lag (async/semi-sync makes RPO/RTO hard to satisfy simultaneously). Cassandra's masterless P2P architecture, append-only commitlog + memtable write path, and native multi-DC replication via NetworkTopologyStrategy with LOCAL_QUORUM solve these.

When Not to Use Cassandra

Avoid Cassandra for complex JOINs, aggregations, fuzzy search, strong transactions (LWT only gives single-partition atomicity at high cost), read-heavy workloads under 1TB (PostgreSQL/MySQL cheaper), and frequent full scans on non-primary-key dimensions (use ClickHouse/Elasticsearch). Choose Cassandra when write QPS > 50k, data > 5TB growing linearly, multi-region active-active with near-zero RPO in local DC, and query patterns are fixed (primary key + time range).

Data Model Decides Everything

Primary Key = Partition Key + Clustering Key

CREATE TABLE device_metric (
  device_id text,
  bucket int,          -- yyyyMMdd, daily bucket
  metric_time timestamp,
  metric_name text,
  value double,
  PRIMARY KEY ((device_id, bucket), metric_time)
) WITH CLUSTERING ORDER BY (metric_time DESC);

The partition key (device_id, bucket) is hashed via Murmur3 to a token range, placing all rows of a partition on the same node (and its replicas). The clustering key metric_time determines on-disk sort order inside the partition and enables efficient range scans. Golden rule: query conditions must resolve to a finite set of partitions ; missing partition key = full table scan.

Data Distribution

Cassandra uses a consistent-hash ring with virtual nodes (default num_tokens=256), greatly reducing skew during scaling compared to manual initial_token. Replication is defined by NetworkTopologyStrategy:

ALTER KEYSPACE metrics_ks WITH replication = {
  'class': 'NetworkTopologyStrategy',
  'dc-shanghai': 3,
  'dc-beijing': 3
};

DC names must match snitch configuration (e.g., GossipingPropertyFileSnitch). Replicas spread across racks so a single rack loss doesn't cause unavailability.

Wide Rows and Sorting

Rows are sparse wide columns; a partition can hold millions of clustering rows. Sort order is fixed at DDL time via WITH CLUSTERING ORDER BY (metric_time DESC) — sorting capability is bought at table creation, not query time , unlike SQL databases.

TTL

TTL is per-cell, set at write: INSERT ... USING TTL 2592000 (30 days). Spring Data supports @TTL annotation or InsertOptions.ttl(). Expired data becomes tombstones, cleaned only during compaction; heavy TTL usage requires compaction strategy adjustments.

Write Path

Write request → sequential append to commitlog (crash recovery) → memtable → flush to immutable SSTables → background compaction merges SSTables and purges expired/duplicate data. No random reads, no in-place updates — UPDATE/DELETE are just new timestamped records. This append-only nature is the root of write performance.

Spring Boot Integration

Dependencies and Configuration

<dependency>
  <groupId>org.springframework.boot</groupId>
  <artifactId>spring-boot-starter-data-cassandra</artifactId>
</dependency>
spring:
  cassandra:
    keyspace-name: metrics_ks
    contact-points: 10.20.0.11,10.20.0.12,10.20.0.13
    port: 9042
    local-datacenter: dc-shanghai
    schema-action: none  # production must disable auto-DDL
    compression: lz4
    request:
      consistency: LOCAL_QUORUM
      timeout: 5s
    connection:
      connect-timeout: 5s
      init-query-timeout: 10s

Critical pitfall: Since Spring Boot 2.6, config prefix changed from spring.data.cassandra.* to spring.cassandra.*; 3.0 removed the old prefix entirely. Many outdated tutorials still use the old prefix, causing silent config failures. Set schema-action: none in production; manage DDL via Liquibase-Cassandra or CI cqlsh scripts.

Driver Customization

application.yml

covers limited options; retry policy, connection pool, load balancing need programmatic DriverConfigLoader:

@Configuration
public class CassandraConfig {
  @Bean
  public CqlSessionBuilderCustomizer sessionCustomizer() {
    return builder -> builder.withConfigLoader(
      DriverConfigLoader.programmaticBuilder()
        .withString(DefaultDriverOption.REQUEST_CONSISTENCY, "LOCAL_QUORUM")
        .withDuration(DefaultDriverOption.REQUEST_TIMEOUT, Duration.ofSeconds(5))
        .withString(DefaultDriverOption.LOAD_BALANCING_LOCAL_DATACENTER, "dc-shanghai")
        .withInt(DefaultDriverOption.CONNECTION_POOL_LOCAL_SIZE, 8)
        .withInt(DefaultDriverOption.CONNECTION_MAX_REQUESTS, 1024)
        .withInt(DefaultDriverOption.RETRY_POLICY_MAX_RETRIES, 2)
        .build()
    );
  }
}
LOAD_BALANCING_LOCAL_DATACENTER

is mandatory; otherwise driver may route to remote DC, doubling latency.

Entity Mapping

@PrimaryKeyClass
public class DeviceMetricKey implements Serializable {
  @PrimaryKeyColumn(name = "device_id", type = PrimaryKeyType.PARTITIONED, ordinal = 0)
  private String deviceId;
  @PrimaryKeyColumn(name = "bucket", type = PrimaryKeyType.PARTITIONED, ordinal = 1)
  private int bucket;
  @PrimaryKeyColumn(name = "metric_time", type = PrimaryKeyType.CLUSTERED, ordinal = 2, ordering = Ordering.DESCENDING)
  private Instant metricTime;
  // getters/setters/equals/hashCode
}

@Table("device_metric")
public class DeviceMetric {
  @PrimaryKey
  private DeviceMetricKey key;
  @Column("metric_name")
  private String metricName;
  @Column("value")
  private double value;
  @Column("tags")
  private Map<String, String> tags;
  @TTL
  private Integer ttl; // seconds, null = no expiry
}
Instant

maps to timestamp, Map<String,String> to map<text,text>. For precise types (uuid, blob, decimal) use @CassandraType(type = Name.UUID) explicitly.

Repository vs CqlSession

public interface DeviceMetricRepository extends CassandraRepository<DeviceMetric, DeviceMetricKey> {
  List<DeviceMetric> findByKeyDeviceIdAndKeyBucket(String deviceId, int bucket, Pageable pageable);
}

Table already has CLUSTERING ORDER BY metric_time DESC; adding OrderByKeyMetricTimeDesc in method name is redundant. Simple CRUD via Repository; complex CQL directly via CqlSession:

@Repository
@RequiredArgsConstructor
public class MetricQueryDao {
  private final CqlSession session;
  public List<DeviceMetric> latest(String deviceId, int bucket, int limit) {
    SimpleStatement stmt = SimpleStatement.builder(
      "SELECT device_id, bucket, metric_time, metric_name, value " +
      "FROM device_metric WHERE device_id = ? AND bucket = ? LIMIT ?")
      .addPositionalValues(deviceId, bucket, limit)
      .setConsistencyLevel(ConsistencyLevel.LOCAL_ONE)
      .setPageSize(500)
      .build();
    return session.execute(stmt).all().stream()
      .map(DeviceMetricMapper::fromRow)
      .toList();
  }
}
all()

fetches all pages; setPageSize only controls network round-trip batch size. For true streaming of large results, use resultSet.iterator() or hardcode LIMIT. SELECT only needed columns — Cassandra is column-oriented; fewer columns reduce network/deserialization overhead, especially on wide tables.

Query Patterns That Don't Blow Up

Avoid ALLOW FILTERING

-- DANGEROUS
SELECT * FROM device_metric WHERE metric_name = 'cpu' ALLOW FILTERING;

This forces coordinator to fan out to all nodes, scan every partition, filter in memory. Works on small data; at tens of millions it stalls the cluster. Only acceptable for rare ad-hoc ops with timeout + LIMIT. Enforce via CI/lint rule: reject any Repository method generating ALLOW FILTERING.

Secondary Indexes and SAI

Traditional 2i are local indexes; equality queries broadcast to all nodes — cost grows with node count. Only for low-frequency, low-write columns. SASI is experimental and disabled by default in 4.0. Recommended: Cassandra 5.0 GA Storage Attached Index (SAI) supporting equality, range, multi-column:

CREATE INDEX idx_metric_name ON device_metric (metric_name) USING 'sai';

Still, indexes maintain local data structures; higher write frequency on indexed columns increases maintenance cost. Indexes supplement, never replace correct partition design.

Materialized Views

MVs auto-maintain a copy with different partition key:

CREATE MATERIALIZED VIEW device_metric_by_name AS
SELECT * FROM device_metric
WHERE metric_name IS NOT NULL AND device_id IS NOT NULL
  AND bucket IS NOT NULL AND metric_time IS NOT NULL
PRIMARY KEY ((metric_name, bucket), metric_time, device_id);

Cost: every base write syncs to MV (write amplification × several), plus consistency window between base and MV. Preference order: redesign data model > application dual-write two tables > materialized view .

Pagination Is Cursor-Based, Not OFFSET

ByteBuffer pagingState = null;
do {
  SimpleStatement stmt = SimpleStatement.newInstance(
    "SELECT * FROM device_metric WHERE device_id = ? AND bucket = ?", devId, bucket)
    .setPageSize(1000)
    .setPagingState(pagingState);
  ResultSet rs = session.execute(stmt);
  rs.forEach(row -> process(row));
  pagingState = rs.getExecutionInfo().getPagingState();
} while (pagingState != null && needMore);

Paging state binds to coordinator; topology change may invalidate it. Business must tolerate "restart from beginning"; caching state across requests in session increases risk.

Batches Only Inside a Single Partition

// Same-partition batch = atomic, correct
BatchStatement batch = BatchStatement.newInstance(BatchType.UNLOGGED);
batch.add(SimpleStatement.newInstance("INSERT INTO device_metric (...) VALUES (...)", ...));

Cross-partition LOGGED batch is an anti-pattern: coordinator writes batchlog, then fans out statements — slower than individual writes, partial rollback on failure. Rule: all statements in a batch must share the partition key . Default batch_size_fail_threshold_in_kb = 50KB; exceed and it errors.

One Table Per Query

Cassandra is query-driven; denormalization is default, not compromise:

Need A: device's daily metrics      → device_metric_by_device
Need B: metric snapshot across devices → device_metric_by_name
Need C: tenant's recent alarms      → alarm_by_tenant

Writer writes multiple tables; each reader hits its own table. Storage redundancy buys zero JOINs, zero network hops — worth it for high-write scenarios.

Consistency Levels and the Math Behind Them

What Each Level Means

ANY : one node (even hinted handoff) acks. Too weak, rarely used.

ONE / LOCAL_ONE : one replica acks. Lowest latency. For telemetry/logs where losing a few is fine. LOCAL_ONE stays in local DC.

QUORUM : majority of all DC replicas (RF=3 → 2). Cross-DC writes wait for remote majority, latency tied to slowest DC.

LOCAL_QUORUM : majority in local DC only. Standard for multi-DC active-active.

EACH_QUORUM : majority in every DC. Latency = slowest DC. Rarely used.

ALL : all replicas. One node down = write failure. Ban in production.

SERIAL / LOCAL_SERIAL : LWT-only (Paxos). Not for regular reads/writes.

The R + W > RF Line

To guarantee read sees latest write: read replicas + write replicas > total replicas. RF=3, QUORUM write + QUORUM read (2+2>3) gives strong consistency. But cross-DC global QUORUM makes write latency depend on slowest DC — unacceptable for multi-active. Hence multi-DC typically uses LOCAL_QUORUM for local consistency, accepts a small bounded inconsistency window across DCs, relies on async replication + background repair.

Read Repair and Anti-Entropy

Two read repair forms: (1) Blocking read repair on QUORUM reads — coordinator compares timestamps, returns latest, schedules async fix for stale replicas. Part of read latency, cannot disable. (2) Probabilistic read repair (old read_repair_chance) removed in 3.0 due to unpredictable write amplification. Now only table option read_repair = BLOCKING (default) or NONE.

More critical: scheduled nodetool repair -pr metrics_ks ( -pr repairs only primary ranges on this node, spreading load). Watch gc_grace_seconds (default 864000s = 10 days) — tombstone retention period. Hard rule: gc_grace_seconds ≥ 2 × repair interval , else deleted data can resurrect.

Lightweight Transactions (LWT): Usable But Not Everywhere

LWT uses Paxos for linearizable single-partition writes: IF NOT EXISTS or IF <condition>.

-- Idempotent dedup
INSERT INTO order_dedup (order_id, created_at) VALUES (?, ?) IF NOT EXISTS;
-- Optimistic lock inventory
UPDATE inventory SET amount = amount - 1 WHERE sku = ? IF amount = ?;

Java checks [applied]:

public boolean decrement(String sku, long expected) {
  SimpleStatement stmt = SimpleStatement.newInstance(
    "UPDATE inventory SET amount = amount - 1 WHERE sku = ? IF amount = ?",
    sku, expected)
    .setConsistencyLevel(ConsistencyLevel.LOCAL_QUORUM)
    .setSerialConsistencyLevel(ConsistencyLevel.LOCAL_SERIAL); // MUST set
  Row row = session.execute(stmt).one();
  return row != null && row.getBoolean("[applied]");
}
setSerialConsistencyLevel

is separate from normal consistency; omitting it breaks semantics. Costs: (1) 4×+ latency vs normal write; (2) single-partition only; (3) high contention → WriteTimeoutException; (4) [applied]=false returns current row, requiring app-level conflict handling. Use only for idempotent dedup, inventory decrement, distributed locks — low-frequency, must-be-correct ops. High-throughput paths: normal write + timestamp (LWW) + app-level idempotency key.

Performance Tuning Knobs

Compaction Strategy

STCS (default): write-optimized, low write amplification, high space amplification, reads may scan many files.

LCS : read-optimized, low space amplification, write amplification ~10×. Good for read-heavy, frequently updated tables.

TWCS for time-series: merges SSTables by time window, drops entire window on expiry, clean tombstone removal.

ALTER TABLE device_metric WITH compaction = {
  'class': 'TimeWindowCompactionStrategy',
  'compaction_window_unit': 'DAYS',
  'compaction_window_size': 1
} AND gc_grace_seconds = 864000;

Paired with daily bucket partition key, one window = one batch of complete partitions; expired window dropped entirely — tombstones largely avoided . Ideal for time-series.

Compression

ALTER TABLE device_metric WITH compression = {
  'class': 'ZstdCompressor',
  'chunk_length_in_kb': 16,
  'level': 3
};

LZ4 default: fast, ~2× ratio. Zstd: 3–5× ratio, higher CPU. Worth it when disk tight or data cold. Time-series + TWCS + Zstd currently optimal combo.

Caching

Key cache (partition index positions): keep enabled, size 5–10% of heap.

Row cache (full rows): disable — competes with memtable for heap, low hit rate on hot updating rows.

Real cache: OS page cache. Give most RAM to OS (e.g., 32GB machine → 8GB heap, 24GB page cache). SSTable reads served from file cache beats JVM caches.

JVM and GC

# cassandra-env.sh
JVM_OPTS="$JVM_OPTS -XX:+UseG1GC"
JVM_OPTS="$JVM_OPTS -XX:MaxGCPauseMillis=300"
JVM_OPTS="$JVM_OPTS -XX:+ParallelRefProcEnabled"

Heap ≤ 8GB. 32GB RAM → 8GB heap, rest page cache. Larger heap = longer Full GC pauses; if Stop-the-world > request timeout, clients error. G1 default since 4.0; 5.0 on JDK 17 offers ZGC for tail latency but G1 remains safer default.

Tombstone Management

Tombstones from deletes load into memory during scans; > tombstone_failure_threshold (default 100k) aborts query. Sources: (1) many single-row DELETEs → batch delete per partition or use TTL; (2) writing null → use CQL unset (omit column) instead; (3) oversized partition key range → bucket by day/month, keep partition ≤ 100MB (prefer ≤ 10MB); (4) full collection overwrites → partial update UPDATE ... SET tags['k'] = v; (5) TTL data not compacted in time → TWCS with aligned compaction_window_size fixes.

Monitor via nodetool tablestats metrics_ks.device_metric — track avg/max tombstones per read (last 5 min) and partition count; alert on upward trend.

Time-Series Modeling Essentials

Bucketing mandatory: bucket = yyyyMMdd (or month) — never let partition grow unbounded.

Partition size ≤ 100MB, ideally ≤ 10MB — eases repair/streaming.

Clustering key ordered by most common range filter.

Cell-level TTL for auto-expiry — cleaner than explicit deletes.

TWCS window granularity must align with bucket granularity to enable whole-window drops.

Production Pitfalls That Actually Happen

Hot partitions: low-cardinality partition key (status, type) concentrates writes. Prefix with high-cardinality field ( device_id) or composite ( (tenant_id, bucket)).

ALLOW FILTERING avalanche: one bad full-scan query exhausts coordinator thread pool, cascades timeouts, cluster-wide jitter.

Cross-partition batch: looks efficient, actually slower; avoid.

Cross-partition LWT: unsupported; attempts yield timeouts and partial success — reconciliation nightmare.

Driver pool misconfig: each app instance opens CONNECTION_POOL_LOCAL_SIZE (default 1, set 4–8) connections per node; CONNECTION_MAX_REQUESTS (default 1024) queues beyond that. Total connections across instances must not saturate node. Driver 4.x merged TokenAwarePolicy into DefaultLoadBalancingPolicy — just configure local DC, don't wrap.

Oversized statements: partition ≤ 100MB, batch ≤ 50KB. Large blobs → object store, keep reference in Cassandra.

Clock drift: client time skew breaks LWW. Enforce NTP. Manual USING TIMESTAMP extremely risky — wrong timestamp makes data permanently invisible.

Forget nodetool cleanup after scaling: old nodes retain migrated data, wasting disk. Run on every old node , not just new ones.

Quick Comparison

MongoDB: replica set + sharding, flexible schema, rich secondary indexes, multi-doc transactions (4.0+). Higher ops complexity when sharded; multi-active DIY. Better for evolving schema, flexible queries, limited ops capacity.

HBase: HDFS + ZooKeeper dependency. Linear write scale, strong Hadoop/Spark integration for hybrid offline/online. Multi-DC via HBase Replication — more complex than Cassandra. Stay if already deep in Hadoop ecosystem with mature ops team.

MySQL sharding: no stack change, full ACID/JOIN/complex queries. Works up to hundreds of GB / few TB with ShardingSphere. Real pain: resharding (16→32 shards = weeks of dual-write/sync/verify), multi-DC nearly impossible (usually one-way sync).

Cassandra: query flexibility weakest (must follow primary key); transactions only single-partition LWT; rich consistency dial (ONE→ALL); ops complexity moderate (no external deps but tuning has learning curve).

Final Word

Cassandra's steep part is almost entirely in modeling. Get the model right, code becomes trivial — no JOINs, no complex transactions, just direct primary-key access. Get it wrong, and it punishes you with timeouts, hot nodes, tombstone storms.

Before going live, repeatedly verify: list every query before DDL, one table per query; partition key determines max scale — bucketing unavoidable for high write; consistency level is latency-vs-correctness knob — multi-DC default LOCAL_QUORUM, reserve LWT for the few ops that truly need it. Do these, Cassandra runs stably for years. Don't, and it will show you exactly where the design broke.

Original Source

Signed-in readers can open the original source through BestHub's protected redirect.

Sign in to view source
Republication Notice

This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactadmin@besthub.devand we will review it promptly.

CompactionData ModelingSpring BootTime SeriesTombstoneCassandraConsistency LevelsSAIHigh Write ThroughputLWT
Xiaolin Talks Programming
Written by

Xiaolin Talks Programming

Focuses on sharing original technical insights. Senior architect at a top tech company with years of experience in technical architecture and management, and extensive interview experience. Offers one-on-one technical coaching, guiding you from beginner to architecture design to technical management. Follow for free learning resources. Free one-on-one interview coaching to help you land offers quickly.

0 followers
Reader feedback

How this landed with the community

Sign in to like

Rate this article

Was this worth your time?

Sign in to rate
Discussion

0 Comments

Thoughtful readers leave field notes, pushback, and hard-won operational detail here.