Hyper-Converged Database (HCD) architecture

Hyper-Converged Database (HCD) is a self-managed, distributed, NoSQL database built on Apache Cassandra® that empowers you to manage your data infrastructure with enterprise-grade capabilities. It provides all the capabilities of Apache Cassandra as well as support for vector search.

HCD is designed to handle large scale data workloads across multiple nodes with no single point of failure. This architecture is based on the understanding that system and hardware failures can and do occur.

Because HCD is based on Apache Cassandra, many of the core architecture concepts and essential distributed database mechanisms are inherited from Cassandra.

Write consistency

HCD addresses the problem of failures by employing a peer-to-peer distributed system across homogeneous nodes where data is distributed among all nodes in the cluster. Each node frequently exchanges state information about itself and other nodes across the cluster using the peer-to-peer gossip protocol.

Each node maintains a sequentially written commit log that captures write activity to ensure data durability. The node then indexes and writes the data to an in-memory structure called a memtable that resembles a write-back cache. When the memtable is full, the node writes the data to disk in an immutable SSTable data file. HCD automatically partitions and replicates all writes throughout the cluster.

HCD periodically consolidates SSTables using compaction. This process consolidates SSTables and discards obsolete data that has been marked for deletion with a tombstone. Deletes are temporarily marked with tombstones to allow time to propagate the deletion across the cluster. To ensure all data across the cluster stays consistent, HCD provides various automatic and manual repair mechanisms.

For more information, see the following:

Request handling

The HCD architecture allows authorized roles to connect to any node in any datacenter to manage the database. Client-side access uses the Cassandra Query Language (CQL) API. CQL uses a similar syntax to SQL and works with table data. You can access CQL through the CQL shell and Apache Cassandra drivers.

Your client applications can send read and write requests to any node in the cluster. When a client connects to a node with a request, that node serves as the coordinator node for that particular client operation. The coordinator acts as a proxy between the client application and the nodes that store the data being requested. Distribution of data across nodes is determined by the replication strategy and partitioning.

For more information, see the following:

Fundamental structures

The following components are essential to understanding the HCD architecture:

Cluster

A group of distributed nodes for storing data. A cluster can have a one or more nodes and one or more single datacenters.

Node (replica)

Nodes store your data. They are the basic database infrastructure component. Technically, they are Java virtual machines running an instance of the HCD software. The primary configuration file for each node is cassandra.yaml.

In addition to storing data and serving read/write requests, nodes can have additional temporary or permanent responsibilites within the cluster:

  • Coordinator node: A temporary designation for a node that receives a request from a client, and then coordinates the request across the cluster.

  • Seed node: A node that helps new nodes join a cluster by providing initial contact points for the gossip protocol. The process of a new node joining a cluster is called bootstrapping. The list of seed nodes is configured per node in the cassandra.yaml file.

  • Replica node: A node that stores a copy of the data for a particular partition. The number of replica nodes for a particular partition is determined by the replication factor of the keyspace that contains the data. A group of replica nodes for the same data can be referred to as a replica set.

Datacenter

A group of related nodes configured together within a cluster for replication purposes. A datacenter can be a physical datacenter or virtual datacenter. Using separate datacenters prevents transactions from being impacted by other workloads and lowers latency. Depending on the replication factor, data can be written to multiple datacenters. Datacenters must never span physical locations.

Schema (keyspaces, tables, primary keys)

The schema defines how data is organized and related within the database. These definitions have a direct impact on how data is stored, accessed, and distributed across the cluster.

Commit logs, memtables, and SSTables

For durability, the database handles writes in stages. Initially, writes are captured in memory (memtables) and a durable commit log. If a node goes down or restarts, the commit log can be replayed to recover lost in-memory writes.

Memtables are periodically flushed to disk in the form of immutable sorted-string tables (SSTables). SSTables are append only, stored on disk sequentially, and maintained for each database table. Compaction periodically consolidates and rewrites SSTables to optimize storage and improve read performance.

Tombstone

A marker in a row that indicates deleted data. Tombstones are propagated to all replicas before the compaction process permanently drops the data from the SSTables. Tombstone management is an important factor for database performance.

Partitioner

Partitioning is important for load balancing in a distributed database. The partitioner determines which replica receives the first copy of a piece of data, and how to allocate data across related replicas. Technically, a partitioner is a hash function that derives a token from the primary key of a row of data. For this reason, the primary key is a determining factor in data distribution across replicas; poorly structured primary keys can lead to uneven data distribution and hotspots.

Replication

Replication is an essential component of distributed databases that ensures data availability by maintaining multiple copies of data on multiple nodes across the cluster. The number of copies is set by the replication factor in the keyspace definition. In production deployments, the replication factor is typically 3 or more, but it cannot exceed the total number of nodes in the cluster.

Snitch

A snitch maps from the IP addresses of nodes to physical and virtual locations, such as racks and datacenters. Snitches inform the database about the network topology for efficient request routing and replica placement.

You must configure a snitch when you create a cluster. There are several snitch types available depending on your deployment environment and topology. All snitches use a dynamic snitch layer that monitors performance and chooses the best replica for reads. The dynamic snitch is enabled by default and recommended for use in most deployments.

Gossip

A peer-to-peer communication protocol to discover and share location and state information about the other nodes in the cluster. Nodes retain gossip information through restarts and outages, and they attempt to reconnect with the cluster using the last known gossip information.

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