Apache Cassandra Archives - A Bias For Action
Musings of a developer
Sun, 19 Jun 2022 22:19:50 +0000
en-AU
hourly
1
CAP Theorem
#respond
Mon, 28 Aug 2017 05:15:39 +0000
The CAP theorem is a tool used to makes system designers aware of trade-offs while designing networked shared-data systems. CAP has influenced the design of many distributed data systems. It made designers aware of a wide range of tradeoff to consider while designing distributed data systems. Over the year the CAP theorem has been widely […]
The post CAP Theorem appeared first on A Bias For Action.
]]>
The CAP theorem is a tool used to makes system designers aware of trade-offs while designing networked shared-data systems. CAP has influenced the design of many distributed data systems. It made designers aware of a wide range of tradeoff to consider while designing distributed data systems. Over the year the CAP theorem has been widely misunderstood tool used to categorize databases. There is much misinformation floating around CAP. Most blog posts around CAP are historical and possibly incorrect.
It is important to understand CAP so that you can identify a lot of the misinformation around it.
The CAP theorem applies to distributed systems that stores state. Eric Brewer at the 2000 Symposium on Principles of Distributed Computing (PODC) conjectured that in any networked shared-data system there is a fundamental trade-off between consistency, availability, and partition tolerance. In 2002 Seth Gilbert and Nancy Lynch of MIT published a formal proof of Brewer’s conjecture.
The theorem states that networked shared-data systems can only guarantee/strongly support two of the following three properties:
Consistency – A guarantee that every node in a distributed cluster returns the same, most recent, successful write. Consistency refers to every client having the same view of the data. There are various types of consistency models. Consistency in CAP (used to prove the theorem) refers to linearizability or sequential consistency a very strong form of consistency.
Availability – Every non-failing node returns a response for all read and write requests in a reasonable amount of time. The key word here is every. To be available every node on (either side of a network partition) must be able to respond in a reasonable amount of time.
Partition Tolerant – The system continues to function and uphold its consistency guarantees in spite of network partitions. Network partitions are a fact of life. Distributed systems guaranteeing partition tolerance can gracefully recover from partitions once the partition heals.
The C and A in ACID represent different concepts than C and in A in the CAP theorem.
The CAP theorem categories systems into three categories:
CP (Consistent and Partition Tolerant) – At first glance, the CP category is confusing, i.e., a system that is consistent and partition tolerant but never available. CP is referring to a category of systems where availability is sacrificed only in the case of a network partition.
CA (Consistent and Available) – CA systems are consistent and available systems in the absence of any network partition. Often a single node DB servers are categorized as CA systems. Single node DB servers do not need to deal with partition tolerance and are thus considered CA systems. The only hole in this theory is that single node DB systems are not a network shared data system and thus do not fall under the preview of CAP. [^11]
AP (Available and Partition Tolerant) – These are systems that are available and partition tolerant but cannot guarantee consistency.
A Venn diagram or a triangle is frequently used to visualize the CAP theorem. Systems fall into the three categories that depicted using the intersecting circles.
CAP Theorem
The part where all three sections intersect is white because it is impossible to have all three properties in networked shared-data systems.
A Venn diagram or a triangle is an incorrect visualization of the CAP. Any CAP theorem visualization such as a triangle or a Venn diagram is a misleading.
The correct way to think about CAP is that in case of a network partition (a rare occurrence) one needs to choose between availability and partition tolerance. Instead of choosing two is more like choose one.
In any networked shared-data systems partition tolerance is a must. Network partitions, dropped messages are a fact of life and must be handled appropriately. Consequently, system designers must choose between consistency and availability. Simplistically speaking a network partition forces designers to either choose perfect consistency or perfect availability. Picking consistency means not being able to answer a clients query as the system cannot guarantee to return the most recent write. This sacrifices availability. Network partition force nonfailing node to reject clients request as these nodes cannot guarantee consistent data. At the opposite end of the spectrum being available means being able to respond to a clients request but the system cannot guarantee consistency, i.e., the most recent value written. Available systems provide the best possible answer under the given circumstance.
During normal operation (lack on network partition) the CAP theorem does not impose constraints on availability or consistency.
The CAP theorem is criticized for being too simplistic and often misleading [^10] [^11]. A decade after the release of the CAP theorem Brewer acknowledge that the CAP theorem oversimplified the choices available in the event of a network partition. According to Brewer, the CAP theorem prohibits only a “tiny part of the design space: perfect availability and consistency in the presence of partitions, which are rare”. System designers have a broad range of options for dealing and recovering from network partitions. The goal of every system must be to “maximize combinations of consistency and availability that make sense for the specific application”
The CAP theorem is a simple straw man to make system designers aware of trade-offs while designing networked shared-data systems. It is a simple starting point and has been widely used to design and discuss tradeoff in NoSQL database.
References:
Brewer’s conjecture and the feasibility of consistent, available, partition-tolerant web services
CAP Twelve Years Later: How the “Rules” Have Changed
Please stop calling databases CP or AP
The post CAP Theorem appeared first on A Bias For Action.
]]>
feed/
0
Reasons for unbalanced Cassandra Cluster
/unbalanced-cassandra-cluster/
/unbalanced-cassandra-cluster/#comments
Fri, 16 Jun 2017 00:52:35 +0000
/?p=380
Sometimes an Apache Cassandra cluster can end up in an unbalanced state. An unbalanced state is where data is unevenly distributed across a cluster or locally configured data directories. There are a number of reasons this can happen. In this blog post, I will cover two basics reasons this might happen. A cluster can end up […]
The post Reasons for unbalanced Cassandra Cluster appeared first on A Bias For Action.
]]>
Sometimes an Apache Cassandra cluster can end up in an unbalanced state. An unbalanced state is where data is unevenly distributed across a cluster or locally configured data directories. There are a number of reasons this can happen. In this blog post, I will cover two basics reasons this might happen. A cluster can end up in an unbalanced state due to two basic reasons:
Configuration
Data Model and data distribution
Configuration
The main configuration option that affects data distribution is the num_token value set in the cassandra.yaml. By default, this is set to 256 but can be configured prior to bootstrapping a node. The num_token configuration is useful when configuring a non-uniform Cassandra cluster. This enables users to distribute data according to the capacity of each node.
Data Model and Data Distribution
Data can also end up unevenly distributed if your data model is not designed for your data profile. You must have a good understanding of your domain and your data model must be fit for your domain. In Cassandra, your partition key is responsible for your distributing data across a cluster and must be carefully chosen.
Example
The best way to understand the above two points is via an example. In this example, we are going to start off with a single node cluster. This node is configured with a num_token value of 256. It has two data directories (/var/lib/cassandra/data, /var/lib/cassandra/data1) specified in the data_file_directories property.
In this example, we will
Bootstrap a single node
Add data to this node.
Add a second node to the cluster and examine how this influences data distribution.
To carry out the test I have created the following schema.CREATE KEYSPACE distribution_test WITH REPLICATION = { 'class' : 'SimpleStrategy' , 'replication_factor' : 1 };
CREATE TABLE dd_test (
partition_key text,
cluster_key text,
test_data text,
PRIMARY KEY (partition_key, cluster_key)
);The distribution_test keyspace has a replication factor 1 and uses the simple replication strategy. The distribution_test keyspace has a table called dd_test. This table has a composite primary key. The primary key is composed of a partition key and cluster key aptly named partition_key and cluster_key respectively.
I have added data to the dd_test table using the following Python script.from cassandra.cluster import Cluster
from concurrent.futures import ThreadPoolExecutor
from logging import thread
import random
import string
cluster = Cluster()
session = cluster.connect('distribution_test')
def insert(from_num, to_num):
for num in range(from_num,to_num):
long_string = ''.join(random.SystemRandom().choice(string.ascii_uppercase + string.digits) for _ in range(50)) * 1000
session.execute("INSERT INTO dd_test (partition_key, cluster_key, test_data) VALUES ('partition_key_{}', 'cluster_key_example_{}' , '{}')".format(num, long_string, long_string))
print "Executed {}".format(num)
futures = []
executor = ThreadPoolExecutor(max_workers=20)
for num in xrange(0,100000,10000):
futures.append(executor.submit(insert, num, num+4999))
for future in futures:
future.result()Notice line 13. The insert statement ensures that the partition key is changed on every insert. As we will see later in the post that this is key to ensuring even data distribution.
After inserting 100000 records I ran nodetool status and got the following values:
As expected all the data is owned by a single node. I also proceeded to check data distribution across the two data directories. Data is evenly spread out. Check out the screen shot below.
Bootstrapping A New Cassandra Node
Let’s add a new node to this cluster. When adding a new node ensure that the auto_bootstrap property is set to true. This ensures that new node automatically joins the new cluster and there is no need for manual configuration. On bootstrapping node2 I see the following log output:INFO [main] 2017-06-01 07:19:45,199 StorageService.java:1435 - JOINING: waiting for schema information to complete
INFO [main] 2017-06-01 07:19:45,250 StorageService.java:1435 - JOINING: schema complete, ready to bootstrap
INFO [main] 2017-06-01 07:19:45,251 StorageService.java:1435 - JOINING: waiting for pending range calculation
INFO [main] 2017-06-01 07:19:45,251 StorageService.java:1435 - JOINING: calculation complete, ready to bootstrap
INFO [main] 2017-06-01 07:19:45,251 StorageService.java:1435 - JOINING: getting bootstrap token
INFO [main] 2017-06-01 07:19:45,341 StorageService.java:1435 - JOINING: sleeping 30000 ms for pending range setup
INFO [main] 2017-06-01 07:20:15,342 StorageService.java:1435 - JOINING: Starting to bootstrap...
INFO [main] 2017-06-01 07:20:15,562 StreamResultFuture.java:90 - [Stream #c219d430-469a-11e7-8af3-81773e2a69ae] Executing streaming plan for Bootstrap
INFO [StreamConnectionEstablisher:1] 2017-06-01 07:20:15,568 StreamSession.java:266 - [Stream #c219d430-469a-11e7-8af3-81773e2a69ae] Starting streaming to /172.20.0.3
INFO [StreamConnectionEstablisher:1] 2017-06-01 07:20:15,591 StreamCoordinator.java:264 - [Stream #c219d430-469a-11e7-8af3-81773e2a69ae, ID#0] Beginning stream session with /172.20.0.3
INFO [STREAM-IN-/172.20.0.3:7000] 2017-06-01 07:20:16,369 StreamResultFuture.java:173 - [Stream #c219d430-469a-11e7-8af3-81773e2a69ae ID#0] Prepare completed. Receiving 4 files(4.046MiB), sending 0 files(0.000KiB)
INFO [StreamReceiveTask:1] 2017-06-01 07:20:32,489 StreamResultFuture.java:187 - [Stream #c219d430-469a-11e7-8af3-81773e2a69ae] Session with /172.20.0.3 is complete
INFO [StreamReceiveTask:1] 2017-06-01 07:20:32,529 StreamResultFuture.java:219 - [Stream #c219d430-469a-11e7-8af3-81773e2a69ae] All sessions completed
INFO [StreamReceiveTask:1] 2017-06-01 07:20:32,535 StorageService.java:1491 - Bootstrap completed! for the tokens [2170857698622202367, 8517072504425343717, -3254771037524187900, -1597835042001935502, 6878847904605480741, 6816916215341820068, 6291189887494617640, -6855333019196580358, -6353317035112065873, 8838974905234547016, 8539810544447438397, -2357950949959511387, 1242077960532340887, 2914039668080386735, 3548015300105653368, 8973388453035242795, -2325235809399362967, -7078812010537656277, 768585495224455336, 7153512700965912517, 8625819392009074153, -6138302849441936958, -2594051958993427953, -735827743339795655, -8202727571538912843, 2180751358288507888, -7872842094207074012, -2926504780761300623, -3197260822146229664, 3411052656191450941, -9049284186987733291, -157882351668930258, 454637839762305232, -2305675997627138050, 5785282040753174988, 8604531769609599767, 4363117061247143957, -7255854383313210529, -3497611663121502480, -6788457421774336480, -7809767930173770420, 6591540654522244365, 1773283733607350132, 1347769111173669066, -7242556233623424655, -1552552727731631642, -1226243976028310059, -8221762326275074149, -7963893314043006091, -850542197910474448, 4219437099703910566, -8039365343972054221, 7756456412568178996, 4057327843751741693, 7155628666873897485, -483058846775660782, 6968839681845709305, 6396337738827005745, -5285173481531605912, 7254663657455123842, 871654822989271789, -604574593420741277, -2446461444470484127, -3707613591745746278, -26727542030118959, -7190990795521107837, 5388348291571480415, 4249499356533972018, 82469082189512791, -6389351372873749061, 5138413916027470955, 2542233707258091740, -4057927973990056143, 552933169018893618, -8237860380097407047, 6917383508758068288, 543382311932406672, -5671560690999322491, -1240369858424929757, 7394536427227616773, 4716882285905136652, 8260705434779371419, 3259812719139852593, -73864539388331289, -3573980475038135246, -1047139059901238511, -1734886021153324482, 8674873751672827600, 3564384074427511950, 2754071903665103098, -1230493021099846761, -2731315467436512731, -7845984767828231726, -8082165594257396645, -2298177264815779081, -3645421111048544165, 9142633389925493379, 7206663288804675578, 2305939212045070856, -5101738026249032246, 6268847697773786891, 5903922100677671597, -2001787466557152206, 1318502870562311928, 5784020265166141829, 5385229217299505171, 6010414616247875068, -8080602674779008196, -9189764569651551963, -8969124116887255329, -9040482343274988119, -8575947267671214955, -1786409930636352174, -757203989676123224, -6640569567328853730, 8431839804447545665, 6781635966829972979, -8328382509754233304, -3181089993114819214, 3243262023331941781, 4213737472390389773, -4046361821170607634, 8877904009116429296, -6931048276693039052, 4838006612846181604, -5561480934050473057, -470112649587309682, 3175935810873308999, -1693695808908080717, -3753035103371291265, -2607412695999984337, -8454963020263227780, 2037428931895594762, 1158209127301347406, -8092787384269386871, -7741092217712244823, 3213269181965324853, 3972662756438857798, -1808499161350392500, -5552429155141285488, 2768019514490102470, -2381885168558935166, 7598271141891576988, -5968675860637356104, -2178161882622813874, -2782395662355757709, 4662660828871465894, 6726990970215445064, 691223925843765893, -3732536320705038428, 2169053732177722520, -3467691490997179851, -3201672755574011994, -706634586120752453, -7297234535099792750, 4195063085031070570, 1024797669903232596, -5102042883065498245, -4295412307398568491, -686079656172478689, -7652228004418103329, -5922734755429174917, 6562442130946482224, 3419893918407781185, -7840156446781283061, -3209525913297552052, -8254134338430272746, -7543559272928655856, -6413145215334169356, 8189387753488304279, 133576451402117013, 6859840908124654784, -480832477584575919, -4466949465409090307, 2224334850433074431, -1077941859365184025, 7694746877711316656, -1541425238506019058, -2694798376156556512, -752352477219592169, 1911593773128947549, 1053380063512932771, 2369074212175473237, 1511764544820953277, -916955813829019462, -3389255868702958135, -4853221732440365560, 436528405098383237, 1482375829041075317, -5200032867873743382, -2359116065936986284, 2323558012249686328, -8714456372962073495, 8264655970455461946, -6377550880012897671, 1904431293553435268, 4281876827769946039, 129146370702108712, -7734952979230886745, -8242033642929677698, 2814257954427384288, -3054904268940900861, -450367636631970004, 8109127225549270088, 8020800599127010921, -2679593348517582224, -1752859515012359396, 4479567094782500113, 6631436750259638425, -5598101192868431054, 2794504876633839794, -1756288804973602742, -534202438069841833, 7027493063976551471, 3982737257386418900, -930887821667260064, 273533908252754639, -1440739443874625754, -2507143240715618713, 2102474087594804091, -7697404848019572605, -8658346826860002424, -4596011304911479578, 6459167476271897881, 6805481881299404326, -7080476038981754572, 5767237378480651194, -978605465701070447, -2322705636227730109, -681083377956364436, -5892502600903861513, -3301298109594166057, 8752787692307586763, -7189731757057511702, 8976826014517372173, -6975924228178934185, 4829963829825594588, -2886015764818641366, -7753659165428942667, -5923615337534732994, 7275900434506691605, 4164818915231369384, 5154263935677458810, 301044058442370317, -7764756071469383634, 7327453618420416121, -2412303784295409238, 781599758351322350, 6625371277047208299, -7611076121968222555, 6127681154442724891, 7313542022544470677, 7222174446262831041, -9138442343184845349, -4110447459792797189, 2534623752979121051]
INFO [main] 2017-06-01 07:20:32,544 StorageService.java:1435 - JOINING: Finish joining ringThe auto bootstrap process:
Added the node to the configured cluster.
Calculated its token range and sets up the new token ranges across the two nodes.
Streamed data from node1 to node 2. This is data that use to belong to node one but now belongs to node two.
After bootstrapping the node I ran nodetool status and got the following output:
Note that node1(172.20.0.3) still has 64.18Mb of data while node two (172.20.0.4) has 31MB of data. This is a bit odd as the new node should have half of the cluster’s data. The reason for the imbalance is because node1 not only contains data the is responsible for but also contains data that is now stored in node2. Running a nodetool cleanup fixes this issue as shown in the screenshot be below:
How has the second node affected the spreading of the data across the two data directories in node1. As expected the data in the two data directories on node1 are evenly spread but half of what they use to be.
All in all, this is a happy situation as our Apache Cassandra cluster is in a balanced state.
Let’s bootstrap node two again but this time configured with a num_token of 10:
Notice that node1 holds over 95% of the data. Note this is expected due to our num_token configuration.
Data Modelling Errors
This is the most frequent reason for an unbalanced cluster. We are going to insert data into the dd_test table but this time we are only going to insert data into a single partition. The code used to insert this data can be found below. The only difference from the code above is line 13.from cassandra.cluster import Cluster
from concurrent.futures import ThreadPoolExecutor
from logging import thread
import random
import string
cluster = Cluster()
session = cluster.connect('distribution_test')
def insert(from_num, to_num):
for num in range(from_num,to_num):
long_string = ''.join(random.SystemRandom().choice(string.ascii_uppercase + string.digits) for _ in range(50)) * 1000
session.execute("INSERT INTO dd_test (partition_key, cluster_key, test_data) VALUES ('partition_key', 'cluster_key_example_{}' , '{}')".format(long_string, long_string))
print "Executed {}".format(num)
futures = []
executor = ThreadPoolExecutor(max_workers=20)
for num in xrange(0,100000,10000):
futures.append(executor.submit(insert, num, num+4999))
for future in futures:
future.result()The above program ends up inserting data into a single partition. One running node tool status we see the following:
We have inserted 485.63 MB of data into node one. Let’s examine how this data is distributed around our two disks.
Note that data1 directory has a majority of data. This is because data is distributed across disks by evenly splitting the nodes token ranges across disks. Since all our data has been inserted into a single partition we have ended up with all our data in a data disk.
Let’s add a new node to the cluster and see how this influences data distribution.
Note adding a new node has resulted in highly imbalanced data distributions. This is again because we have inserted all data into a single partition and thus resulted in a highly unbalanced node.
I hope this has helped you understand why we might end up with unbalanced Cassandra nodes.
Moral of the story. Make sure you put in effort into designing your schema correctly.
The post Reasons for unbalanced Cassandra Cluster appeared first on A Bias For Action.
]]>
/unbalanced-cassandra-cluster/feed/
2
Understanding an Apache Cassandra Memtable Flush
/apache-cassandra-memtable-flush/
/apache-cassandra-memtable-flush/#comments
Tue, 30 May 2017 00:15:15 +0000
/?p=361
A recent question in the Apache Cassandra mailing list triggered this blog post. The question revolved around events that trigger a memtable flush. Understanding the root cause of a memtable flush is essential to get a better understanding of Apache Cassandra. Another question that frequently crops up is the size of an SSTable as a result of […]
The post Understanding an Apache Cassandra Memtable Flush appeared first on A Bias For Action.
]]>
A recent question in the Apache Cassandra mailing list triggered this blog post. The question revolved around events that trigger a memtable flush. Understanding the root cause of a memtable flush is essential to get a better understanding of Apache Cassandra. Another question that frequently crops up is the size of an SSTable as a result of a memtable flush.
A memtable is created for every table or column family. There can be multiple memtables for a table but only one of them will be active. The rest will be waiting to be flushed. There are a few properties that affect a memtables size and flushing frequency. These include:
memtable_flush_writers – This is the number of threads allocated for flushing memtables to disk. This defaults to two.
memtable_heap_space_in_mb – This is the total allocated space for all memtables on an Apache Cassandra node. By default, this is one-fourth your heap size. Specifying this property results in an absolute heap size in MB as opposed to a percentage of the total JVM heap.
memtable_cleanup_threshold – A percentage of your total available memtable space that will trigger a memtable cleanup. memtable_cleanup_threshold defaults to 1 / (memtable_flush_writers + 1). By default this is essentially 33% of your memtable_heap_space_in_mb. A scheduled cleanup results in flushing of the table/column family that occupies the largest portion of memtable space. This keeps happening till your available memtable memory drops below the cleanup threshold.
commitlog_total_space_in_mb – The total space in MB that is reserved for the commit log. If unspecified this defaults to the smaller of the two numbers i.e. 8192 MB or 25% of the total space of the commit log volume.
memtable_flush_period_in_ms – This is a CQL table property that specifies the number of milliseconds after which a memtable should be flushed. This property is specified on table creation.
memtable_allocation_type – Stipulates where Cassandra allocates and manages memtable memory. Memory can be allocated on the JVM heap or directly into memory. This does not affect the flushing of memtables but only signifies where memtables space is allocated.
So when does a memtable get flushed to disk? A memtable is flushed when:
The commit log reaches its maximum size – The main aim of the commit log is to track all data that has not been written to disk i.e. data in a memtable which has not been flushed to disk. Commit log forces flushing of memtables if a commit log runs out of disk space. The commit log allocates chunks of space in what is called a commit log segment. When a commit log runs out of disk space (surpasses its config threshold) it needs to recycle allocated segements. This recycling process triggers flushing of all memtables. The commit log cannot be cleared as long as it refers to data in a memtable. Doing so risks data loss. Thus when a commit log is full all memetables are first flushed to disk and then the commit log is recycled. The org.apache.cassandra.db.commitlog.AbstractCommitLogSegmentManager.maybeFlushToReclaim() method is where this recycling is done.
Periodically – If the CQL table sets the memtable_flush_period_in_ms property then the memtable gets flushed after the configured number of milliseconds has elapsed.
Memtable surpasses it on and off-heap memory threshold – This is best understood using an example. Let assume we have an Apache Cassandra instance that has allocated 4G of space. Out of this only 3,925.5MB is available to the Java runtime. Please look at the following StackOverflow question for the reasons behind this. Of this, by default, we have 981 MB allocated towards memtable i.e. 1/4the of 3,925.5. Our memtable_cleanup_threshold is the default value i.e. 33 percent of the total memtable heap and off heap memory. In our example that comes to 327 MB. Thus when total space allocated for all memtables is greater than 327 MB a memtable clean-up is triggered. The cleanup process looks for the largest memtable and flushes that to disk. The org.apache.cassandra.db.ColumnFamilyStore.FlushLargestColumnFamily class can be examined for further details.
Due to the above configuration options and varying Apache Cassandra workloads, our SSTable size on disk can vary greatly. One thing to remember is that by default SSTables are compressed. SSTable compression can be turned off using compression table property.
The post Understanding an Apache Cassandra Memtable Flush appeared first on A Bias For Action.
]]>
/apache-cassandra-memtable-flush/feed/
6
Cassandra Query Language (CQL) Tutorial
/cassandra-query-language-cql-tutorial/
/cassandra-query-language-cql-tutorial/#comments
Fri, 26 May 2017 22:28:20 +0000
/?p=308
Apache Cassandra and the Cassandra Query Language (CQL) have evolved over the past couple of years. Key improvements include: Significant storage engine improvements Introduction of SSTable Attached Secondary Index i.e SASI Indexes Materialized views Simple role based authentication This post is an updated to “A Practical Introduction to Cassandra Query Language”. The tutorial will concentrate […]
The post Cassandra Query Language (CQL) Tutorial appeared first on A Bias For Action.
]]>
Apache Cassandra and the Cassandra Query Language (CQL) have evolved over the past couple of years. Key improvements include:
Significant storage engine improvements
Introduction of SSTable Attached Secondary Index i.e SASI Indexes
Materialized views
Simple role based authentication
This post is an updated to “A Practical Introduction to Cassandra Query Language”. The tutorial will concentrate on two things:
Cassandra Query Language and its interaction with the new storage engine.
Introducing the various CQL statements via a practical example.
CQL Overview
Cassandra Query Language or CQL is a data management language akin to SQL. CQL is a simple API over Cassandra’s internal storage structures. Apache Cassandra introduced CQL in version 0.8. Thrift an RPC-based API was the preferred data management language prior to the introduction of CQL. CQL provides a flatter learning curve in comparison to Thrift and thus its is the preferred way of interacting with Apache Cassandra. In fact, Thrift support is slated to be removed in Apache Cassandra 4.0 release.
CQL has many restrictions in comparison to SQL. These restrictions prevent inefficient querying across a distributed database. CQL queries should not visit a large number of nodes to retrieve required data. This has the potential to impact cluster-wide performance. Thus CQL prevents the following:
No arbitrary WHERE clause – Apache Cassandra prevents arbitrary predicates in a WHERE statement. Where clauses must have columns specified in your primary key.
No JOINS – You cannot join data from two Apache Cassandra tables.
No arbitrary GROUP BY – GROUP BY can only be applied to a partition or cluster column. Apache Cassandra 3.10 added GROUP BY support to SELECT statements.
No arbitrary ORDER BY clauses – Order by can only be applied to a clustered column.
CQL Fundamentals
Let’s start off by understanding some basic concepts i.e. cluster, keyspaces, tables aka column family and primary key.
Apache Cassandra Cluster – A cluster is a group of computers working together that are viewed as a single system. A distributed database is a database system that is spread across a cluster. Apache Cassandra is a distributed database spread across a cluster of nodes. Think of Apache Cassandra Cluster as a database server spread across a number of machines.
Keyspaces – A keyspace is similar to an RDBMS schema/database. A keyspace is a logical grouping of Apache Cassandra tables. Like a database, a keyspace has attributes that define system-wide behaviour. Two key attributes are the replication factor and the replication strategy. The replication factor defines the number of replicas the replication strategy defines the algorithm used to determine the placement of the replicas.
Tables – An Apache Cassandra table is similar to an RDBMS table. Like in an RDBMS a table is made up of a number of rows. This is where the similarity ends. Think of an Apache Cassandra table as a map of sorted maps. A table contains rows each of which is accessible by the partition key. The row contains column data which are ordered by the clustering key. Visualise an Apache Cassandra table as a map of sorted maps spread across a cluster of nodes.
Primary Key – A Primary key uniquely identifies an Apache Cassandra row. A primary key can be a simple key or a composite key. A composite key is made up of two parts, a partition key and a cluster key. The partition key determines data distribution in the cluster while the cluster key determines sort order within a partition.
The diagram below helps visualise the above concepts.
Apache Cassandra Cluster, Keyspace, Table and Primary Key Overview
Notice how keyspaces and tables are spread across the cluster. In summary, Apache Cassandra cluster contains keyspaces. Keyspaces contain tables. Tables contain rows which are retrieved via their primary key. The goal is to distribute this data across a cluster of nodes.
Installing Cassandra
To go through this tutorial you need Apache Cassandra installed. You can either use a Docker base Apache Cassandra installation or use any of the installation methods laid out in 5 ways to install Apache Cassandra. I will be using the Docker base Cassandra installation.
CQLSH
Cqlsh is a Python based utility that enables you to execute CQL. A cqlsh prompt can be obtained by simply executing the cqlsh utility. In the Docker-based installation navigate to the containers command prompt. Once at the command prompt simply type cqlsh. For native installs the cqlsh utility can be found in CASSANDRA_INSTALLTATION_DIR/bin. Executing the cqlsh utility will take you to the cqlsh prompt.
CQL By Examples
Navigating an Apache Cassandra Cluster
The first things you need to get familiar with is navigating around a cluster. The DESCRIBE command, or DESC or shorthand, is the key to navigating around the cluster. The DESCRIBE command can be used to list cluster objects such as keyspaces, tables, types, functions, aggregates. It can also be used to output CQL commands to recreate a keyspace, table, index, materialized view, schema, function, aggregate or any CQL object. To get a full list of DESCRIBE command options please type HELP DESC at the cqlsh prompt. The HELP provided is self-explanatory.
Describe command help
Often the first thing you want to do when you are the CQLSH prompt is to list keyspaces. TheDESCRIBE KEYSPACEScommand provides a list of all keyspace. Go ahead and execute this command. Since we have freshly installed Apache Cassandra all you will see a list of system keyspaces. Systems keyspaces help Apache Cassandra to keep track of cluster and node related metadata.
Here are a couple of gotchas. There is no command to list all indexes. To list indexes you have to query the system keyspace. To list all indexed you will need to execute the following command:SELECT * FROM system.”IndexInfo”;Similarly, there is no command to list all materialized views. To get a list of all the views you must execute the following command:SELECT * FROM system_schema.views;
CQL Create Keyspace
Let’s create our first keyspace. The keyspace is called animals. Execute the statement below to create the animal keyspace.CREATE KEYSPACE animals WITH REPLICATION = { 'class' : 'SimpleStrategy' , 'replication_factor' : 3 };Two key parameters to take note of is the class and the replication factor. The class defines the replication strategy i.e. algorithm used to determine the placement of the replicas. The replication factor determines the number of replicas. Since I have three node cluster I have chosen a replication factor of 3.
Execute theDESC KEYSPACEScommand once more to see if the aminal keyspace is created. Next, let’s connect to the created keyspace with the help the of the USE command. The USE command enables switching context between keyspaces. Once a keyspace is chosen all subsequent commands are executed in the context of the chosen keyspace. Please executeUSE animals;
CQL Create Table
Let’s create a table called monkeys in the animals keyspace. Execute the following command to create the monkey table.CREATE TABLE monkeys (
type text,
family text,
common_name text,
conservation_status text,
avg_size_in_grams int,
PRIMARY KEY ((type, family), common_name)
);Take special note of the primary key. The primary key defined above is a composite key. The primary key has two parts. i.e. partition key and cluster key. The first column of the primary key is your partition key. The remaining columns are used to determine the cluster key. A composite partition key, a partition key made up of multiple columns, is defined by using an extra set of parentheses before the clustering columns. The partition key helps distribute data across the cluster while the cluster key determines the order of the data stored within a row. Thus type and family is our composite partition key and common_name is our cluster key. When designing a table think of the partition key as a tool to spread data evenly across a cluster while the cluster key helps determine the order of that data within a partition. Your data and query patterns will influence your primary key. Please note the cluster key is optional.
CQL INSERT
Let’s insert some data into above table. When inserting data the primary key is mandatory.INSERT INTO monkeys (type, family, common_name, conservation_status)
VALUES ('New World Monkey', 'Cebidae', 'white-headed capuchin', 'Least concern');Execute a select statement to see that the data has been inserted successfully.SELECT * FROM monkeys;You should see the following output:
Note columns in red are your partition key. Cyan coloured columns are your cluster key. The columns in purple are the rest of your columns.
CQL Consistency Level
Cassandra enables users to configure the number of replicas in a cluster that must acknowledge a read or write operation before considering the operation successful. The consistency level is a required parameter in any read and write operation and determines the exact number of nodes that must successfully complete the operation before considering the operation successful.
Let’s check the consistency level set in our cqlsh client. To check your consistency level simply executeCONSISTENCYYou should seeCurrent consistency level is ONE.This implies that only one node needs to insert data successfully for the statement to be considered successful. The statement is executed on all replicas but only needs to be successfully written to one in order to be successful. On each node, the insert statement is first written to the commit log and then into a memtable i.e. a write back cache. The memtable is flushed to disk either periodically or when a certain size threshold is reached. The process of flushing converts a memtable into an SSTable.
Nodetool Flush
Let’s manually flush the above insert to disk so that we can examine the created SSTable. You can force a flush by executing the following at the command prompt.nodetool flush animalsThe above command flushes all tables in the animals keyspace to disk. Navigate to the animal keyspace data directory to view the results. Look at your YAML file for your data directory location. In the Docker-based installation the directory is:/var/lib/cassandra/data/animals/monkeys-generated_hex_stringThe flushed SSTable can be found in a file with the suffix -Data.db. My file is called mc-1-big-Data.db. We can output the contents of the file using a command line tool sstabledump. The utility sstabledump outputs the contents of -Data.db file as JSON. To get the JSON output simply executesstabledump mc-1-big-Data.dbMy mc-1-big-Data.db has the following contents[
{
"partition" : {
"key" : [ "New World Monkey", "Cebidae" ],
"position" : 0
},
"rows" : [
{
"type" : "row",
"position" : 43,
"clustering" : [ "white-headed capuchin" ],
"liveness_info" : { "tstamp" : "2017-05-18T22:55:47.631382Z" },
"cells" : [
{ "name" : "conservation_status", "value" : "Least concern" }
]
}
]
}
]Note the partition object contains the partition keys used while the rows array contains the row data for the partition. The rows object contains a clustering array that stores our cluster column related data. The cells array in the row object contains all additional column data.
Let’s insert some additional data to see how this affects the underlying storage.INSERT INTO monkeys (type, family, common_name, conservation_status, avg_size_in_grams)
VALUES ('New World Monkey', 'Cebidae', 'white-fronted capuchin', 'Least concern', 3400);
INSERT INTO monkeys (type, family, common_name, conservation_status, avg_size_in_grams)
VALUES ('New World Monkey', 'Cebidae', 'white-headed capuchin', 'Least concern', 3900);
INSERT INTO monkeys (type, family, common_name, conservation_status, avg_size_in_grams)
VALUES ('New World Monkey', 'Cebidae', 'tufted capuchin', 'Least concern', 4800);
INSERT INTO monkeys (type, family, common_name, conservation_status, avg_size_in_grams)
VALUES ('New World Monkey', 'Cebidae', 'blond capuchin', 'Least concern', 2800);
INSERT INTO monkeys (type, family, common_name, conservation_status)
VALUES ('Old World Monkey', 'Colobinae', 'Nilgiri langur', 'vulnerable');Flush memtable data to disk by running the following command:nodetool flush animalsThe first thing you should notice is a new file mc-2-big-Data.db. Note mc-1-big-Data.db has not been updated to as SSTables are immutable.
Please run the following command to output JSON output for mc-2-big-Data.db.sstabledump mc-2-big-Data.dbYou should get the following output.[
{
"partition" : {
"key" : [ "Old World Monkey", "Colobinae" ],
"position" : 0
},
"rows" : [
{
"type" : "row",
"position" : 45,
"clustering" : [ "Nilgiri langur" ],
"liveness_info" : { "tstamp" : "2017-05-19T02:19:12.639248Z" },
"cells" : [
{ "name" : "conservation_status", "value" : "vulnerable" }
]
}
]
},
{
"partition" : {
"key" : [ "New World Monkey", "Cebidae" ],
"position" : 82
},
"rows" : [
{
"type" : "row",
"position" : 125,
"clustering" : [ "blond capuchin" ],
"liveness_info" : { "tstamp" : "2017-05-19T02:18:37.670932Z" },
"cells" : [
{ "name" : "avg_size_in_grams", "value" : "2800" },
{ "name" : "conservation_status", "value" : "Least concern" }
]
},
{
"type" : "row",
"position" : 167,
"clustering" : [ "tufted capuchin" ],
"liveness_info" : { "tstamp" : "2017-05-19T02:18:37.662353Z" },
"cells" : [
{ "name" : "avg_size_in_grams", "value" : "4800" },
{ "name" : "conservation_status", "value" : "Least concern" }
]
},
{
"type" : "row",
"position" : 210,
"clustering" : [ "white-fronted capuchin" ],
"liveness_info" : { "tstamp" : "2017-05-19T02:18:37.643547Z" },
"cells" : [
{ "name" : "avg_size_in_grams", "value" : "3400" },
{ "name" : "conservation_status", "value" : "Least concern" }
]
},
{
"type" : "row",
"position" : 258,
"clustering" : [ "white-headed capuchin" ],
"liveness_info" : { "tstamp" : "2017-05-19T02:18:37.657137Z" },
"cells" : [
{ "name" : "avg_size_in_grams", "value" : "3900" },
{ "name" : "conservation_status", "value" : "Least concern" }
]
}
]
}
]Note we now have data in two partition keys. Observe that in the second partition row data is sorted by the cluster key. Notice that the second insert did not error out even thought that data already existed. This is because inserts in Cassandra are actually upserts. An upsert inserts the row if it does not exist otherwise updates the existing data.
CQL DELETE
Let’s delete some data. Please execute the following delete command:DELETE FROM monkeys WHERE type='New World Monkey' and family='Cebidae' and common_name='white-headed capuchin'Again runnodetool flush animalsIn your data directory, you should see a new file mc-3-big-Data.db.
Please run the following command to convert mc-3-big-Data.db into JSON.sstabledump mc-2-big-Data.dbYou should see the following output:[
{
"partition" : {
"key" : [ "New World Monkey", "Cebidae" ],
"position" : 0
},
"rows" : [
{
"type" : "row",
"position" : 43,
"clustering" : [ "white-headed capuchin" ],
"deletion_info" : { "marked_deleted" : "2017-05-19T05:05:23.822752Z", "local_delete_time" : "2017-05-19T05:05:23Z" },
"cells" : [ ]
}
]
}
]The first thing to notice is that deletion does not result in any actual data being physically deleted. This is because SSTables are immutable and are never modified. In fact, the delete statement caused an insertion of a special value called a tombstone. A tombstone records deletion related information. There are many kinds of tombstones. The delete statement above caused the insertion of a row tombstone.
CQL UPDATE
Like inserts updates in Cassandra are an upsert. The following is an example of an update statement.UPDATE monkeys
SET conservation_status = 'least concern' , avg_size_in_grams = 3000
WHERE type='Old World Monkey'
AND family='Colobinae' AND common_name='Southern plains gray langur'Even though the primary key specified in the update statement does not exist data will still be inserted. This is due to the upsert semantics of a CQL update statement. In the case, the primary key already exists the appropriate values will be updated. Feel free to flush the animal keyspace and view the data file created.
CQL Time To Live aka TTL
A compelling feature in Cassandra is the ability to expire data. This feature is especially useful when dealing with time series data. The TTL feature enables you to expire columns after a set number of seconds. Below is an insert example with TTLINSERT INTO monkeys (type, family, common_name, conservation_status, avg_size_in_grams)
VALUES ('Old World Monkey', 'Colobinae', 'Tarai gray langur', 'near threatened', 3000) USING TTL 600;Using the TTL function you can query the number of seconds left for the column to expire. Please note a TTL is only assigned to non-primary key columns.SELECT TTL(avg_size_in_grams), TTL(conservation_status) from monkeys;Use the nodetool flush command to check data saved to disk. You should see something similar to the following JSON.{
"partition" : {
"key" : [ "Old World Monkey", "Colobinae" ],
"position" : 0
},
"rows" : [
{
"type" : "row",
"position" : 146,
"clustering" : [ "Tarai gray langur" ],
"liveness_info" : { "tstamp" : "2017-05-20T22:17:19.028766Z", "ttl" : 600, "expires_at" : "2017-05 20T22:27:19Z", "expired" : false },
"cells" : [
{ "name" : "avg_size_in_grams", "value" : "3000" },
{ "name" : "conservation_status", "value" : "near threatened" }
]
}
]
}Look at the liveness_info object. It has an addition of ttl and expires_at element as opposed to previous data file outputs in this tutorial. The expires_at specific the exact UTC time when the columns will expire. The ttl element specifies the TTL value passed in on insert.
When a TTL is present on all colums of a row then on the expiration of the TTL the entire row is omitted from select statement results. When the TTL is only present on particular columns then only that particular columns data is omitted from the results.
CQL Data Types
CQL supports an array of data types which includes character, numeric, collections, and user-defined types. The following table outlines the supported data types.
Category
Data Type
Description
Example
Numeric data type
tinyint
8 bit signed integer. Values can range from −128 to 127.
3
smallint
16 bit signed integer. Values can range from −32,768 to 32,767
20000
int
32-bit signed int. Values can range from −2,147,483,648 to 2,147,483,647, from
3