Create and evaluate data models and schemas

Data models and the resulting schemas (keyspaces, tables, and indexes) influence the performance of your DataStax Enterprise (DSE) cluster and associated applications that you develop. Optimal data modeling typically leads to better performance, scalability, and stability for your DataStax Enterprise (DSE) cluster and applications. Suboptimal data models can lead to poor performance, even with small amounts of data.

This guide introduces important data modeling concepts, strategies you can use to evaluate your data model, and considerations for properly tuned schemas.

Prepare your data model

Data modeling includes theoretical designs as well as the practical outcomes of those designs. For example:

  • Logical and physical models

  • Data access patterns (queries)

  • Keyspaces and replication

  • Indexes

  • Table properties like compaction and compression

  • Table schema elements like columns and primary keys

  • Hot data and caches

To get started with data modeling, see the following:

If you are migrating from a relational database, be aware that there are significant differences between RDBMS and Cassandra-based data models.

Notably, data models are based on query patterns, not entities and relationships. Additionally, distributed databases scale horizontally by adding more nodes to the cluster rather than relying on a single machine.

Often, developers migrating from relational databases are attracted to materialized views, secondary indexes, and user-defined functions (UDFs). However, these patterns aren’t ideal for heavy usage in distributed NoSQL databases due to performance impacts at scale. Review this guide and the related information to understand the purpose and requirements for different data types and auxiliary indexes.

Tune table properties

Reads, writes, and data maintenance are fundamental database operations. In addition to designing an efficient data model, there are many ways to tune these operations for optimal performance and resource usage. To determine the ideal configuration for your workloads, you must understand how these operations work, the relevant settings, and the tradeoffs involved with major tuning choices. For example, workloads that are sensitive to read latency can be tuned for faster reads but this usually comes at a cost of slower writes.

The following sections focus on specific table properties related to reads, writes, and data maintenance, such as compaction and compression. This guide isn’t a comprehensive settings reference or a detailed primer for these concepts. Use this information in conjunction with related documentation. For example:

To check the properties for an existing table, use the DESCRIBE TABLE command.

Compaction strategy

It is important that you choose an appropriate compaction strategy for your use case and data model. Misconfigured or unsuitable compaction strategies can degrade performance and overconsume system resources. The following table summarizes the supported compaction strategies and general configuration guidance for each strategy.

Supported compaction strategies
Strategy Use case Configuration notes General disk requirements

UnifiedCompactionStrategy (UCS)

Unifies and builds on tiered (STCS) and leveled (LCS) compaction.

Recommended for all workloads with the exception of time series data with expiring time-to-live (TTL) workloads that is better suited to TimeWindowCompactionStrategy (TWCS).

Understand how the scaling property works and make sure it is set to your preferred mode and performance targets.

Many compaction properties for this strategy are sufficient at the default values, but you might need to tune them for certain workloads.

Disk space requirements depend on the configured mode (STCS, LCS, balanced, or multiple tier-specific modes).

The maximum disk space for this strategy is set by the max_space_overhead parameter. If the default configuration is 20 percent of your node’s free disk space, DataStax recommends tuning this parameter.

SizeTieredCompactionStrategy (STCS)

Good for write-heavy workloads that prioritize fast writes over read latency. DataStax recommends UCS in tiered mode over traditional STCS.

In most cases, the default properties are sufficient. If this strategy produces too many outliers, or compaction runs more often that you would like, then the size range and thresholds might need tuning.

Make sure there is sufficient memory and storage available for the number and size of SSTables as well as overhead for the compaction process. The sum of all SSTables being compacted must be smaller than the remaining disk space, ideally less than 50 percent.

Avoid exceeding 50 percent of free disk space, which is likely to occur with manual compaction where all SSTables are merged into one giant SSTable.

LeveledCompactionStrategy (LCS)

Good for read-heavy workloads that perform best with fewer SSTables. DataStax recommends UCS in leveled mode or multiple tier-specific modes over traditional LCS, which can alleviate some performance concerns at high levels (L3 and above).

Understand the mechanisms and resource requirements at L0 compared to L1 and higher. Requires tuning memtable parameters for optimal performance, such as less frequent flushing of memtables to avoid overloading L0.

Due to guaranteed non-overlapping row key ranges, LCS requires much less disk space for compaction compared to STCS. However, you must account for the use of STCS as a failsafe at L0, which requires more disk space.

Furthermore, the maximum overhead for LCS increases dramatically beyond L3 because each level is approximately 10 times larger than the preceding level. I/O saturation is possible when compacting at the highest levels due to progressively larger SSTables at each additional level. For more information, see LCS compaction write amplification and disk requirements.

TimeWindowCompactionStrategy (TWCS)

Designed for time series data and expiring time-to-live (TTL) workloads, especially data that is written once, in chronological order, and never updated.

TWCS must be enabled when you create a table. You cannot apply TWCS retroactively to existing tables that weren’t created with the proper time windowing layout.

Similar to STCS, TWCS requires a maximum disk space overhead of 50 percent of the total size of SSTables in the last created bucket. To ensure adequate disk space, determine the size of the largest bucket or window ever generated, and divide that value by 2: TWCS disk space = Largest bucket size / 2. For new deployments, you must monitor and tune this during cluster performance tests.

If node metrics report a high number of pending compactions, you might need to tune settings or resources related to compaction.

Compaction can be tuned by table properties and cassandra.yaml parameters. For more information about compaction properties and tuning, including monitoring and testing compaction strategies, see Configure compaction.

Compression

Table compression maximizes the storage capacity of DSE nodes by reducing the volume of data on disk and disk I/O, particularly for read-heavy workloads. Unlike relational databases, compression in DSE doesn’t degrade write performance because compression happens after writes are processed.

Compression is enabled by default with the LZ4Compressor algorithm if not specified in the CREATE TABLE command. DataStax recommends compression with the default configuration for most tables, particularly with DSE Search.

Avoid increasing the chunk size without proper testing because it can significantly degrade read performance if the database has to read unnecessarily oversized chunks compared to the actual amount of data requested by a read. However, if your data model’s query patterns usually request many rows, a larger chunk size might be more performant.

If you plan to use Transparent Data Encryption (TDE), DataStax recommends that you also enable compression.

For information about when to use compression, compression algorithms, and other compression properties, see Configure compression.

Caches

All databases read fastest when frequently access data is stored in memory. DSE caches data retrieved from reads as well as auxiliary indexes. For more information about caches used for read requests, see Read operations.

You can inspect cache usage with nodetool info. If cache is full and the hit ratio is below 90%, you might need to tune the cache configuration.

Read performance tuning with table properties

Depending on your workloads, you might need to tune the following table properties for more ideal read performance:

Get the schema

After creating your keyspaces and tables in your DSE cluster, get the schema to evaluate your data model.

The schema is included in a diagnostic tarball generated by DSE OpsCenter or the DataStax diagnostic collection scripts.

You can also get the schema with cqlsh:

  1. Connect to your DSE cluster using cqlsh.

  2. Run the DESCRIBE SCHEMA command and pipe the output to a file:

    cqlsh -e 'DESCRIBE SCHEMA;' > schema.cql

The examples in this guide assume you have a running DSE cluster with keyspaces and tables already created, and you have output the schema definition to a file named schema.cql. If you use a different file name, edit the example commands accordingly.

For diagnostic tarballs, schema.cql is located at /driver/schema for each node.

The schema is retrieved from a single node. Typically, this is acceptable because schema changes are propagated to all nodes in the cluster.

In situations where you need to compare nodes or investigate desynchronized nodes, use nodetool commands. These commands run on every node in a cluster to provide a complete report of the cluster state. Then, you can compare the output for the anomalous nodes.

Evaluate keyspace replication

Check that all keyspaces have appropriate replication properties:

Always prefer NetworkTopologyStrategy over SimpleStrategy

NetworkTopologyStrategy is required for proper balancing of multi-datacenter clusters and rack-based configurations. If you have only one datacenter, NetworkTopologyStrategy makes it easier to add datacenters in the future.

Failure to set or alter replication strategy when required

When you add or move a datacenter, you must manually reconfigure replication for the system_auth keyspace.

When deploying a node to production, you must set the replication factors for security keyspaces and, if applicable, DSE Analytics keyspaces.

Avoid over-replication, under-replication, and selective replication that skips datacenters

If keyspaces are under-replicated inside datacenters or not replicated to all datacenters, you can lose access to clusters completely. For example, the system_auth keyspace is critical for authentication and authorization. In contrast, excessive replication can overload the cluster with additional internode networking and high I/O.

Replicas aren’t a disaster recovery solution

Don’t use replication as your sole disaster recovery and backup solution. Other ways to protect and duplicate data include adding datacenters and enabling backup retention.

Set the replication factor to an odd number, typically 3 or 5

Replication greater than 5 is typically not necessary.

Even numbers of replicas usually don’t perform as well as odd numbers with consistency levels like QUORUM or LOCAL_QUORUM because even numbers of nodes make the cluster less resilient to failure. For example, the number of nodes required for QUORUM and LOCAL_QUORUM are calculated as N / 2 + 1, where N is the number of replicas in a cluster (for QUORUM) or datacenter (for LOCAL_QUORUM). If the replication factor is 2, the number of replicas in QUORUM is 2, and operations fail if 1 node is down. However, if the replication factor is 3, operations won’t fail with 1 node down because the number of replicas for QUORUM is still 2.

The following table compares the tolerance of various replication factors for down nodes when the consistency level is QUORUM or LOCAL_QUORUM. Specifically, it compares the following:

  • Down node tolerance: Based on the replication factor, this is the number of nodes that can be down without losing the ability to achieve consistency level QUORUM in a cluster or LOCAL_QUORUM in a datacenter.

  • Down node failure: Based on the replication factor, this is the number of down nodes that cause operations to fail due to missed consistency level QUORUM or LOCAL_QUORUM.

Replication factor QUORUM or LOCAL_QUORUM threshold Down node tolerance Down node failure

2

2 / 2 + 1 = 2

0

1

3

3 / 2 + 1 = 2

1

2

4

4 / 2 + 1 = 3

1

2

5

5 / 2 + 1 = 3

2

3

6

6 / 2 + 1 = 4

2

3

7

7 / 2 + 1 = 4

3

4

Fix replication problems

To fix problems with replication, you can either run the ALTER KEYSPACE command manually, or use a script to perform these actions automatically.

  • System keyspaces that use LocalStrategy or EverywhereStrategy must be left unchanged.

  • You must run nodetool repair on all keyspaces where a change is made.

For example, the following script can help you adjust the replication factor of system and user-provided keyspaces based on the current cluster topology. It generates a file with CQL commands that you can run after reviewing them.

adjust-keyspaces.sh script
#!/bin/bash

rep_factor=3
fix_system=1
keyspaces=("system_auth" "system_distributed"  "dse_security" "solr_admin" "dse_perf" "dse_leases" "dse_analytics" '"HiveMetaStore"' "system_traces" "dse_advrep" "dse_system")
cqlsh_options=""
nodetool_options=""

function usage() {
    echo "Usage: $0 [-n replication_factor] [-c cqlsh_options_as_a_string] [-o nodetool_options_as_a_string] [keyspace1 keyspace2 ...]"
    echo "Defaults:"
    echo "   replication_factor = $rep_factor"
    echo "   keyspace1.. = name(s) of keyspace(s) to fix replication factor."
    echo "                 fix all system keyspaces if not specified"
}

while getopts ":hc:o:n:" opt; do
    case $opt in
        n) rep_factor=$OPTARG
           ;;
        c) cqlsh_options=$OPTARG
           ;;
        o) nodetool_options=$OPTARG
           ;;
        h) usage
           exit 0
           ;;
        *) usage
           exit 1
           ;;
    esac
done
shift "$((OPTIND -1))"

if [ -z "$(command -v nodetool)" ]; then
    echo "Can't find nodetool in the PATH"
    exit 1
fi

if [ -z "$(command -v cqlsh)" ]; then
    echo "Can't find cqlsh in the PATH"
    exit 1
fi

if [ $# -ne 0 ] ; then
    keyspaces=()
    fix_system=0
    for ks in "$@" ; do
        keyspaces+=("$ks")
    done
fi

my_pid=$$

NODETOOL_FILE=/tmp/nodetool-out-${my_pid}.txt
SCHEMA_FILE=/tmp/cqlsh-schema-${my_pid}.txt

# File with commands to execute
CQL_FILE=/tmp/fix-keyspaces-${my_pid}.cql
rm -f $CQL_FILE
touch $CQL_FILE

nodetool $nodetool_options status > $NODETOOL_FILE
RES=$?
if [ $RES -ne 0 ] ; then
    echo "Can't execute 'nodetool status'! Exit code: $?"
    exit 1
fi

if grep -e '^[UD][NJLM] ' "$NODETOOL_FILE"|grep -v -e '^UN ' > /dev/null 2>&1 ; then
    echo "Cluster has nodes with non-UN status, can't adjust replication factor!"
    exit 1
fi

cqlsh $cqlsh_options -e 'DESCRIBE FULL SCHEMA;' > $SCHEMA_FILE
RES=$?
if [ $RES -ne 0 ] ; then
    echo "Can't get schema via cqlsh! Exit code: $?"
    exit 1
fi

declare -A all_ks
for i in $(grep 'CREATE KEYSPACE' "$SCHEMA_FILE"|grep -e 'SimpleStrategy\|NetworkTopologyStrategy'|sed -e 's|^CREATE KEYSPACE \([^ ]*\).*$|\1|'); do
    all_ks[$i]=1
done
#echo "All keyspaces=${!all_ks[@]}"

declare -A all_dcs

curr_dc=''
cnt=0
while read -d $'\n' line ; do
#    echo "$line"
    if echo "$line"|grep -e '^Datacenter: '  > /dev/null 2>&1 ; then
        new_dc=$(echo "$line"|sed -e 's|^Datacenter: \([^ ]*\).*$|\1|')
#        echo "new_dc=$new_dc"
        if [ -z "$new_dc" ] ; then
            echo "Can't extract DC from line '$line'"
            exit 1
        fi
        if [ -z "$curr_dc" ] ; then
            curr_dc=$new_dc
        else
            max_rf=$rep_factor
            if [ $rep_factor -gt $cnt ]; then
                max_rf=$cnt
            fi
            echo "$curr_dc has $cnt nodes max RF=$max_rf"
            all_dcs[$curr_dc]=$max_rf
            curr_dc=$new_dc
            cnt=0
        fi
    fi
    if echo "$line"|grep -e '^[UD][NJLM] ' > /dev/null 2>&1 ; then
        ((cnt++))
    fi
done < $NODETOOL_FILE
# push the last DC as well
max_rf=$rep_factor
if [ $rep_factor -gt $cnt ]; then
    max_rf=$cnt
fi
echo "$curr_dc has $cnt nodes max RF=$max_rf"
all_dcs[$curr_dc]=$max_rf

#echo "All DCs=${!all_dcs[@]}"
if [ ${#all_dcs[@]} -eq 0 ]; then
    echo "Can't identify data centers!"
    exit 1
fi

to_repair=()
for i in "${keyspaces[@]}"; do
#    echo "Processing $i"
    if [ -z "${all_ks[$i]}" ]; then
        if [ $fix_system = "0" ]; then
            echo "$i not in the list of existing keyspaces!"
        fi
        continue
    fi
    echo -n "ALTER KEYSPACE $i WITH replication = {'class': 'NetworkTopologyStrategy'" >> $CQL_FILE
    for key in "${!all_dcs[@]}"; do
        echo -n ", '$key': ${all_dcs[$key]}" >> $CQL_FILE
    done
    echo "};" >> $CQL_FILE
    to_repair+=("$i")
done

ret_code=0
if [ ${#to_repair[@]} -eq 0 ]; then
    rm -f $CQL_FILE
    echo "No keyspaces processed!"
    ret_code=1
else
    echo "Please execute command 'cqlsh -f $CQL_FILE $cqlsh_options' to adjust replication factor for keyspaces"
    echo "After that, execute following commands on each node of the cluster:"
    for i in "${to_repair[@]}" ; do
        echo "nodetool $nodetool_options repair -pr $i"
    done
fi

# remove already processed files
rm -f $SCHEMA_FILE $NODETOOL_FILE

exit $ret_code
Usage

The script must be executed on one of the nodes in your cluster because it uses cqlsh and nodetool to find cluster topology and extract a list of keyspaces.

The script accepts several command line parameters, and it supports clusters with no authentication, as well as clusters with non-default configuration:

  • -h: Get usage information.

  • -n replication_factor: Specify the desired replication factor for keyspaces. If this parameter is greater than the number of nodes in a specific datacenter, then the number of nodes is used instead.

  • -c cqlsh_options: All options that must be passed to cqlsh, such as authentication.

  • -o nodetool_options: All options that must be passed to nodetool.

  • Optional list of keyspaces for which adjustment should be done. If not specified, a hardcoded list of system keyspaces is used.

Output

The script generates a file with CQL commands for changing replication factors of keyspaces, and it prints usage instructions, including a nodetool repair command for each keyspace that would be changed. You must run nodetool repair on all keyspaces where a change is made. For example:

SearchAnalytics has 3 nodes max RF=3
Please execute command 'cqlsh -f /tmp/fix-keyspaces-7111.cql' to adjust replication factor for keyspaces
After that, execute following commands on each node of the cluster:
nodetool  repair -pr system_auth
nodetool  repair -pr system_distributed
nodetool  repair -pr dse_security
nodetool  repair -pr dse_perf
nodetool  repair -pr dse_leases
nodetool  repair -pr dse_analytics
nodetool  repair -pr "HiveMetaStore"
nodetool  repair -pr system_traces

The script adjusts replication for the following keyspaces:

  • system_auth

  • system_distributed

  • system_traces

  • dse_security

  • solr_admin

  • dse_perf

  • dse_leases

  • dse_analytics

  • HiveMetaStore

  • dse_advrep

  • dse_system

Evaluate the number of tables

Using your schema.cql file, use the following command to check how many tables and keyspaces are in your cluster:

grep 'CREATE TABLE' schema.cql |wc -l

Having too many tables in a cluster can cause substantial performance degradation. Inherently, throughput drops as the number of tables increases.

Because performance is influenced by a range of other factors, it is difficult to determine a universal table threshold. You can use the following guidelines as a starting point:

  • Set the tables_warn_threshold guardrail to 500 or less, depending on your cluster resources and workloads.

    Clusters can support more tables, but this warning alerts you that you might need to audit underutilized tables or start planning to rearchitect your data model and schemas.

  • Set the tables_failure_threshold guardrail to 1000 or less, depending on your cluster resources and workloads.

    Distributed database clusters almost universally experience performance degradation issues at or before 1000 tables, such as excessive memory usage, long-running compactions, and query failures. While some clusters might remain stable, this many tables is usually evidence of an inefficient data model.

    For example, a Cassandra data modeling best practice is one table per major query pattern; however, if your applications use 1000 or more substantially different query patterns, your query patterns could probably be optimized through prepared statements or indexes.

  • Understand how JVM heap and byte count can reduce a cluster’s table threshold:

    • Each table uses a small amount of memory for metadata.

    • All tables produce memtables and, through compaction, SSTables.

    • The number of columns and rows in a table can increase pressure on memory by storing more auxiliary data, such as bloom filters, caches, and indexes.

    • Each keyspace causes additional overhead in JVM memory. Having many keyspaces can reduce your cluster’s actual table capacity, regardless of the number of tables per keyspace.

  • Tune table properties for appropriate resource usage and performance.

Evaluate keys and partitions

Primary keys, partitions, and clustering columns are fundamental to query performance. For an introduction to these concepts, see the following:

Partition size

For efficient operation, partitions must be sized within certain limits. Partition size can be measured by the number of values in the partition and the partition size on disk.

A cell is a stored value in one column in one row. In addition to the actual data value, each cell has associated metadata, such as timestamp, optional TTL, and additional data for more complex data types, such as collections.

Cassandra has a practical limit of 2 billion (231) cells per partition. It is difficult to exactly predict the required disk space because this depends on the number of rows, columns, primary key columns, and static columns in the table. However, for most use cases, DataStax recommends fewer than 100,000 items per partition and partitions smaller than 100 MB.

The presence of large partitions (greater than 100 MB) is a sign of a suboptimal data model, and they create an additional load on the database. For example:

  • Low cardinality of partition keys.

  • Non-uniform spread of data between partitions.

  • Too many columns and rows in the table, especially when every row contains data for all or most columns.

  • Storing large BLOBs or long text strings in cells.

The performance impacts of large partitions include:

  • Imbalanced partitions cause some nodes to run more operations (hotspots) and trigger compaction more often.

  • More data must be transferred when reading the whole partition.

  • Performance degradation in external systems. For example, a partition is the minimal object mapped into a Apache Spark™ partition, so any imbalance in the database partitions can lead to an imbalance when processing data with Spark.

For more information, see the following:

Get information about partitions

If your tables have large partitions, low cardinality partitions, or other issues with partitions, the only long term solution is to change your data model and recreate the tables with more performant partition keys and clustering columns.

Use the following tools to get information about partitions, including the size and cardinality:

nodetool tablehistograms

Use the nodetool tablehistograms command to find the partition sizes for 75, 95, 99, and 100 percentiles:

nodetool tablehistograms KEYSPACE_NAME TABLE_NAME

If the output reports significantly different Partition Size and Cell Count values for each partition, then it is likely that the partition key values are unevenly distributed across partitions, there are too many columns in the table, or too many elements in non-frozen collection columns.

test/widerows histograms
Percentile  SSTables     Write Latency      Read Latency    Partition Size        Cell Count
                              (micros)          (micros)           (bytes)
50%             1.00            126.93           1310.72            545791               103
75%             1.00            152.32           1310.72            545791               124
95%             1.00            785.94           1310.72            545791               124
98%             1.00           1629.72           1310.72            545791               124
99%             1.00           2346.80           1310.72            545791               124
Min             1.00             73.46           1310.72            454827                87
Max             1.00           2346.80           1572.86            545791               124

Similar information can be obtained from the sstablepartitions command.

nodetool tablestats

Information about maximum partition size is available through nodetool tablestats.

In the output, check the following:

  • Compacted partition maximum bytes: Make sure the values are smaller than the recommended 100 MB.

  • Number of partitions: This is an estimate that can indicate low cardinality partition keys.

dsbulk

You can use DataStax Bulk Loader (DSBulk) to identify partition keys that have the largest number of rows if your cluster has a non-uniform spread of partition key values. The main advantage of DSBulk is that it works with the whole cluster. For example, to find the largest partitions in a table:

dsbulk count -k KEYSPACE_NAME -t TABLE_NAME --log.verbosity 0 --stats.modes partitions
Result
'29' 106 51.71
'96' 99 48.29
sstablemetadata

You can use the sstablemetadata utility with the -s command line parameter to identify the largest partitions in specific SSTables.

sstablemetadata provides information about the largest partitions as both row count and size in bytes. However, it runs on individual SSTable files, so it could be difficult to analyze partitions that are split across multiple files.

sstablemetadata -u -s path_to_file/mc-1-big-Data.db
Result
SSTable: /Users/ott/var/dse-5.1/data/cassandra/data/datastax/vehicle-8311d7d14eb211e8b416e79c15bfa204/mc-1-big
Size: 58378400
Partitions: 800
Tombstones: 0
Cells: 2982161
WidestPartitions:
  [266:20180425] 534
  [202:20180425] 534
  [232:20180425] 534
  [187:20180425] 534
  [228:20180425] 534
LargestPartitions:
  [266:20180425] 134568 (131.4 kB)
  [202:20180425] 134568 (131.4 kB)
  [232:20180425] 134568 (131.4 kB)
  [187:20180425] 134568 (131.4 kB)
  [228:20180425] 134568 (131.4 kB)
...
SELECT DISTINCT

To check for low cardinality partition keys, you can run a SELECT DISTINCT query on the partition key columns, and then check the count in the output for the number of distinct values returned.

SELECT DISTINCT partition_key_list, count(*) FROM table

Evaluate the number of columns per table

DataStax doesn’t recommend creating tables with hundreds or thousands of columns for the following reasons:

  • It is easy to exceed the recommended maximum number of cells per partition and columns per row.

  • Every cell has a storage cost.

  • Range scans perform poorly.

If a table has too many columns, analyze access patterns for that table, and then consider the following solutions:

  • If multiple columns are always read together, combine those columns into one frozen user-defined type (UDT) where all data in the UDT is written as one cell per row.

  • Handle serialization and deserialization of data at the application level, not within the table schema.

  • Where appropriate, store values as BLOBs.

  • Split data into separate tables.

Evaluate data types

DSE provides a rich set of data types that can be used for table columns. Because so many data types exist, it can be easy for new users to choose an incorrect or suboptimal data type. Using appropriate data types helps ensure efficient storage and accurate query results.

The following list describes notable antipatterns to avoid when choosing data types. It doesn’t describe all data types or all possible antipatterns.

Avoid overreliance on the text type

Don’t use the text type when another specialized type is more appropriate. It can use excess storage or produce incorrect query results.

For example, don’t use text to store timestamps. A timestamp encoded as text using ISO-8601 notation occupies 28 bytes, while the specialized timestamp type uses only 8 bytes.

Choose numeric types with appropriate precision

Don’t use numeric types that have larger value ranges than are necessary for the data because this stores values at an unnecessary precision, wasting space on irrelevant values. For example, don’t use bigint (8 bytes) when int (4 bytes) is sufficient. The excess storage is even worse when decimal and varint types are used incorrectly because they don’t have a fixed size; they are sized based on the actual value.

Range queries can fail silently when inappropriate types are used for sorted data

Don’t use string types when numeric or chronological ordering is needed. For example, storing numbers or timestamps as text can cause range queries (greater than, less than, equal to) to fail because the text representation of those values doesn’t support the required ordering for comparison. If a sortable type is used, but it isn’t the correct type for the data, then range queries can run successfully but return incorrect results.

Prefer collections (maps, lists, sets) over UDTs and tuples

To store multiple values in a single column, you can use collections (maps, lists, sets), user-defined types (UDTs), and tuples:

  • Collections are preferred when possible, particularly frozen collections.

    Performance degradation is possible when a database has large collections and non-frozen collections. This includes read amplification, non-idempotent writes, increased storage overhead, and extra read-before-write operations. To learn when to use collections, how to create efficient collections, and how to reduce performance impacts, see Collections in CQL.

  • UDTs are useful because they are customizable. However, they are more complex to manage and prone to becoming overloaded with too many elements. Consider collections before falling back to UDTs. When UDTs are used, use them deliberately and implement organizational policies that reduce the likelihood of uncontrolled alterations to UDTs.

  • DataStax recommends collections and UDTs over tuples. Tuples are always frozen, requiring a full rewrite for each update with explicit null values for unset/empty elements, and they require additional handling during upgrade and migrations to avoid data loss.

Collections, UDTs, and tuples all have practical and literal size limits that must be monitored and preemptively addressed to avoid performance degradation.

Prefer frozen collections when using collections

Collection types (map, list, set) can be either frozen or non-frozen. This is an important distinction that can impact application performance and data model designs.

Counters are non-idempotent and restrict the table schema

The counter type is useful for integers that are modified by increments or decrements. However, because counter operations are non-idempotent, they cannot be retried automatically without risking overcounts or undercounts. Additionally, the counter type places restrictions on the table schema and operations like ALTER TABLE and INSERT. For example, tables with counter columns can have only counter columns and the primary key columns. If you choose to use counter columns, be aware of the limitations, and make sure your application logic accounts for the non-idempotent nature of counter operations. For more information, see CQL data types: Counter.

BLOBs have a practical limit of 2 GB

Storing large BLOBs, such as images or videos, can lead to unbalanced partitions and degraded performance. Instead, for large binary content, store the data in an object store or content delivery network (CDN), and then store the URI in your database.

Evaluate auxiliary indexing

By default, you can query tables using partition key columns, including the primary key and clustering columns. To query on non-partition key columns, use secondary (auxiliary) indexes:

  • Storage-Attached Indexing (SAI) (recommended)

  • Original secondary indexing (2i)

  • SSTable-Attached Secondary Indexing (SASI)

  • DSE Search indexes

  • Materialized views (enables non-primary key queries but isn’t a literal index)

For more information about index types, use cases, and CQL commands for creating and querying indexes, see Index types and use cases and Prepare to use materialized views.

Follow these best practices to ensure your indexes are performant and not creating unnecessary load on the cluster:

Indexing isn’t a replacement for an effective data model

Indexes can enable specific query patterns that aren’t otherwise possible, but they aren’t a replacement for best practices like one table per query. Each index has a maintenance and storage cost in the database, and poorly designed indexes lead to poor query performance.

Avoid superfluous indexes and audit unused indexes

Because the database needs resources to build and maintain secondary indexes, DataStax recommends limiting the number of indexes and removing unused indexes. Using your schema.cql file, use the following command to check how many indexes are in your cluster:

grep 'CREATE INDEX' schema.cql|wc -l
Storage-attached indexing (SAI) is recommended

SAI is recommended for most workloads but these indexes require careful planning and preparation of hardware capacity. For more information, see the following:

Original secondary indexing (2i, index lookups) isn’t recommended

2i isn’t recommended except for edge cases that aren’t addressed by other indexing strategies. This type of indexing is prone to performance degradation and only suitable for specific schemas and query patterns. For more information, see the following:

Don’t use SSTable-attached Secondary Indexes (SASI)

SASI is deprecated in DSE 6.8. If you are migrating or upgrading from an earlier version, remove SASI indexes from your schema and data model. Using your schema.cql file, use the following command to check for SASI indexes in your cluster:

grep 'CREATE CUSTOM INDEX.*SASIIndex' schema.cql|wc -l
DSE Search indexes are only for DSE Search workloads

If you are using DSE Search, you can use DSE Search indexes. To create performant indexes, consider the following:

  • Review the DSE Search documentation so you understand how it works.

  • Limitations are imposed on these indexes by Apache Lucene™, Solr, and DSE Search. Many of these limitations are addressed by SAI.

  • The maximum number of indexes your cluster can sustainably support depends on your version of DSE, index characteristics, and the cluster’s hardware. For example, to build an object inside the DSE Search index, the database must read the corresponding row from the base table, which increases I/O. For more information, see Capacity planning for DSE Search.

  • Some columns aren’t supported, such as counter, frozen map, and static columns.

  • A single DSE Search index on a single node cannot contain more than 2 billion documents (number of objects). It is easy to reach this limit on UDT columns, which are indexed as separate documents.

  • DSE Search reads at consistency level ONE. This level means there can be discrepancies in returned results if the data isn’t repaired regularly

  • Row-level access control (RLAC) isn’t supported.

Using your schema.cql file, use the following command to check for DSE Search indexes in your cluster:

grep 'CREATE CUSTOM INDEX.*Cql3SolrSecondaryIndex' schema.cql|wc -l

To get schema and index configuration details, use the CQL DESCRIBE ACTIVE SEARCH INDEX command.

Understand performance implications and limitations of materialized views

Materialized views (MVs) aren’t a literal index type, but they can be used to enable queries that aren’t possible with a table’s primary key. For information about how MVs work, known limiations, and performance impacts, see Prepare to use materialized views.

MVs are considered experimental and not recommended for production use.

They can be useful for data modeling because you can write data to one table, and then use MVs to test different primary keys, query patterns, and slices of data from that one table.

MVs are usually less performant than duplicate writes to multiple tables that have different primary keys. Although these duplicate writes add complexity to your application logic, this approach has advantages over MVs:

  • Complete flexibility in defining the primary key for each table, whereas MVs must respect the base table’s primary key.

  • Avoids performance impacts, such as extra read-before-write cycles when iterating writes to MVs.

  • Avoids repair issues known to occur with materialized views.

Both tables and MVs consume storage space in the cluster because MV data is stored in each view, similar to an index or table. Each MV increases the total size of the data stored. This can be significant when there are many MVs or the base table is large. Depending on your workloads, separate tables or secondary indexes can be more efficient to store.

If you choose to use MVs, plan your data model to minimize the number of MVs that you create, and learn how to maintain MVs. DataStax recommends using the most recent version of DSE 6.8 (minimum) or 6.9 (preferred) to get the most stable version of materialized views. If you are on an earlier version, upgrading is recommended if your data model relies on MVs. Using your schema.cql file, use the following command to check the number of MVs in your cluster:

grep 'CREATE MATERIALIZED VIEW' schema.cql|wc -l

Non-idempotent and read-before-write operations

In any application, non-idempotent operations are avoided because, by design, every run changes the system state. In the event of a failure, automatic retries risk duplicating the operation, and extra failure handling logic to check if an operation was completed can be imperfect, increase latency, and require read-before-write operations. However, not retrying the operation could also leave the system in an undesirable state.

DataStax recommends that you design your application to use idempotent operations whenever possible. Minimize non-idempotent and expensive read-before-write operations, such as lightweight transactions (LWTs) and certain collection mutations, because they can degrade performance and throughput.

In some cases, these operations are unavoidable. If you must use non-idempotent and read-before-write operations, follow industry best practices for these patterns. For example, make sure your applications have appropriate retry logic and your clusters have adequate resources to handle the additional load.

For more information, see the following:

Tombstones and zombies

It’s important that you understand how tombstones work and their impact on database performance. Specifically, you need to know how to address down nodes to avoid resurrecting missed deletes.

Tombstones are markers of deleted data. In a distributed database, tombstones must be propagated to replicas before the data is actually dropped from the database. If a node is down and misses a delete, the data can be incorrectly restored to the cluster when the node recovers, which is known as a zombie record. Therefore, tombstones must live long enough to be propagated, but they must not live so long that performance degrades from accumulating too many tombstones.

The gc_grace_seconds table property sets the expiration for tombstones. To avoid zombies, you must repair down nodes before gc_grace_seconds expires. Read repair isn’t the same as a anti-entropy repair, and you must not rely on read repair to propagate tombstones. For more information about tombstones, repairing nodes, and the grace period property, see Deletes and tombstones.

Was this helpful?

Give Feedback

How can we improve the documentation?

© Copyright IBM Corporation 2026 | Privacy policy | Terms of use Manage Privacy Choices

Apache, Apache Cassandra, Cassandra, Apache Tomcat, Tomcat, Apache Lucene, Apache Solr, Apache Hadoop, Hadoop, Apache Pulsar, Pulsar, Apache Spark, Spark, Apache TinkerPop, TinkerPop, Apache Kafka and Kafka are either registered trademarks or trademarks of the Apache Software Foundation or its subsidiaries in Canada, the United States and/or other countries. Kubernetes is the registered trademark of the Linux Foundation.

General Inquiries: Contact IBM