Skip to main content
Version: 3.0 (next)

Apache Cassandra Apache Cassandra Integration Guide

Connect to Apache Cassandra to read and write rows from your pipelines. This guide covers connection setup, the six function types, and pipeline integration.

Overview​

Apache Cassandra is a distributed wide-column database built for heavy write loads — sensor readings, event history, machine state. The connector speaks the CQL native protocol (version 4), so it also works with ScyllaDB, DataStax Enterprise and DataStax Astra.

The connector provides:

  • CQL queries with values bound to ? markers, a limit sent to the cluster as the page size, and paging state to resume
  • Key reads of a partition or a row, without writing CQL
  • Row writes of up to 1,000 rows per call, with a time to live and an insert-only mode
  • Deletes of a row or a whole partition, optionally only if it exists
  • Single statements — counter updates, lightweight transactions, schema changes
  • Batches of up to 100 statements: logged, unlogged or counter
  • A consistency level per connection and per function
  • Data-centre-aware, token-aware routing: each request goes straight to a replica of its partition in the local data centre
  • Username and password authentication, TLS and mutual TLS

One connection, one keyspace​

A connection addresses one keyspace. Every function names its own table, so a single connection covers every table in the keyspace; a table in another keyspace is named as keyspace.table.

Tables are shaped by their queries​

A Cassandra table is read by its primary key. The first part, the partition key, decides which nodes hold a row; the rest, the clustering columns, order the rows inside a partition:

CREATE TABLE readings (
device_id text,
ts timestamp,
celsius double,
PRIMARY KEY ((device_id), ts)
) WITH CLUSTERING ORDER BY (ts DESC);

A read names the partition key — WHERE device_id = ? — and may narrow it by the clustering columns in order. A query that does not name the partition key has to visit every node, which Cassandra refuses unless the query says ALLOW FILTERING. Design each table for the reads a pipeline will make.

Names​

Cassandra folds an unquoted name to lower case: a column created as deviceId is stored as deviceid. The connector matches JSON field names to columns the same way — an exact match first, then the one column whose name differs only in case — so {"deviceId": …} writes the deviceid column. A name created in double quotes keeps its case and is matched exactly.

Connection Configuration​

Creating an Apache Cassandra Connection​

Navigate to Connections → New Connection → Apache Cassandra and configure the following.

1. Profile Information​

FieldDefaultDescription
Profile Name-A descriptive name for this connection profile (required, max 100 characters)
Description-Optional description for this connection

2. Cluster and Keyspace (Connection tab)​

FieldDefaultDescription
Contact Points-Comma-separated nodes to connect to first — 10.0.0.11, 10.0.0.12. An entry may carry its own port (10.0.0.11:9142) – required
Port9042CQL native transport port of the contact points that do not name their own
Keyspace-The keyspace this connection addresses – required
Local Data Centre-Send requests to this data centre's nodes, using the others only when it has none left. Leave empty for a single data-centre cluster

The rest of the cluster is discovered from the contact points, so two or three nodes are enough. Prefer IP addresses: a name that resolves to several addresses is dialled once per address.

3. Authentication and TLS (Security tab)​

FieldDefaultDescription
Username-Role name for PasswordAuthenticator. Leave empty when the cluster runs with authentication off (AllowAllAuthenticator, the default)
Password-The role's password. Stored encrypted
Enable TLSoffEncrypt the connection. The nodes must have client_encryption_options enabled
Skip Certificate VerificationoffAccept any server certificate. Only displayed with TLS on. Not for production
TLS Server Name (SNI)-Name to verify the nodes' certificates against. Only displayed with TLS on
CA Certificate-PEM of the CA that signed the nodes' certificates. Leave empty to use the system's trusted CAs. Only displayed with TLS on
Client Certificate / Client Private Key-PEM pair for mutual TLS (require_client_auth on the nodes). Provide both or neither. Stored encrypted. Only displayed with TLS on

4. Advanced (Advanced tab)​

FieldDefaultDescription
ConsistencyLOCAL_QUORUMHow many replicas must answer each request, unless a function chooses its own
Connect Timeout10sBound on connecting, including the keyspace check, and on each health check (1s–300s). Each operation has its own Timeout

Consistency levels​

LevelReplicas that must answer
ONE, TWO, THREEThat many, in any data centre
LOCAL_ONEOne in the local data centre
QUORUMA majority of all replicas
LOCAL_QUORUMA majority in the local data centre — the usual choice
EACH_QUORUMA majority in every data centre (writes)
ALLEvery replica
ANYAny node, even one that only stores a hint (writes)
SERIAL, LOCAL_SERIALRead the latest value a lightweight transaction committed (reads)

A level the cluster cannot meet right now — TWO on a one-node cluster, QUORUM with too many replicas down — fails the request as unavailable. That verdict is transient, so Retry on Fail and store-and-forward retry it.

How the health check works​

Test Connection — and every later health check — opens a session and reads the keyspace's row from system_schema.keyspaces. That proves the contact points answer, the credentials are accepted, and the keyspace exists. A keyspace that does not exist fails the check with its name in the message.

Local development with Docker​

The official image runs a one-node cluster with authentication off. Cap its heap: left to the default, the JVM sizes itself from the host's memory and Docker Desktop kills the container (exit code 137).

docker run -d --name cassandra -p 9042:9042 \
-e MAX_HEAP_SIZE=512M -e HEAP_NEWSIZE=128M cassandra:5.0

A node takes about a minute to start answering. Then create a keyspace:

docker exec -it cassandra cqlsh -e \
"CREATE KEYSPACE plant WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 1}"

Connect with Contact Points 127.0.0.1 and Keyspace plant. If MaestroHub itself runs in Docker, use host.docker.internal instead of 127.0.0.1.

Function Builder​

Creating Apache Cassandra Functions​

Open the connection, go to the Functions tab, and click New Function. Pick one of the six operations.

Rows in and out​

Rows are plain JSON in both directions. On the way in, each value is converted to its column's CQL type:

CQL typeWrite it as
text, varchar, asciia string
int, bigint, smallint, tinyint, counter, varinta whole number, or its text — every digit of a 19-digit bigint is kept
float, double, decimala number, or its text
booleantrue / false
timestampRFC 3339 text (2026-09-30T10:00:00Z) or epoch milliseconds
date"2026-09-30"
time"10:15:30.5"
uuid, timeuuidthe UUID's text
inetthe address's text
blobbase64, or 0x-prefixed hex
duration"1h30m"
list, setan array
mapan object (keys converted to the map's key type)
user-defined typean object of its fields

A value that does not fit its column — "seven" for an int, 4294967296 for an int — is refused before anything is sent, naming the value and the column.

On the way out, rows come back as JSON objects keyed by column name, in the same forms: text for UUIDs, addresses, decimals and wide varints; RFC 3339 for timestamps; YYYY-MM-DD for dates; base64 for blobs; an array for a tuple.

Tuples cannot be bound

The driver cannot bind a tuple column to a ? marker. Write Rows refuses such a row and says so. Write a tuple as a literal in Execute Statement instead: UPDATE readings SET pos = (7, 'line-3') WHERE ….

Query​

Run a CQL SELECT. Only SELECT runs here: a role that may only read through this connector must not be able to write through it.

FieldDefaultDescription
Query-The SELECT, with ? markers for values – required
Values-JSON array of values for the markers, in order — ["press-1", "2026-09-01T00:00:00Z"]
Consistencyconnection'sReplicas that must answer
Limit100Maximum rows (1–10,000). Sent to Cassandra as the page size
Paging State-Resume where a previous run stopped
Timeout30mBound on this operation

Paging​

When more rows may follow, the result carries pagingState. Feed it into the next run's Paging State, with the same query and values, to continue exactly after the last row returned. The limit is enforced by the cluster, never by trimming a larger result.

Read Rows​

Read rows by primary key without writing CQL.

FieldDefaultDescription
Table-readings, or plant.readings – required
Key-JSON object of primary key columns — {"device_id": "press-1"} – required
ColumnsallComma-separated columns to return
Consistency, Limit, Paging State, TimeoutAs in Query

The key must name every partition key column. Clustering columns narrow the read in their declared order: with PRIMARY KEY ((device_id), ts, seq), a key may name ts alone but not seq without ts. A key that matches nothing returns no rows — not an error.

Write Rows​

Write one row or up to 1,000.

FieldDefaultDescription
Table-The table to write – required
Rows-One JSON object, or an array of objects – required
Insert OnlyoffWrite a row only if no row with its primary key exists (IF NOT EXISTS)
Time to Live (seconds)-Seconds until the written values expire (0 or empty keeps them, or uses the table's default_time_to_live)
Consistencyconnection'sReplicas that must acknowledge each write
Timeout30mBound on this operation

Every row is checked against the table — its columns, its primary key, every value's type — before the first row is written, so a mistake in row 900 does not leave 899 rows written. A field the connector does not recognise makes it read the table's columns again once, in case the table was altered.

Cassandra writes are upserts: a row with the same primary key is overwritten. Insert Only keeps the first write instead; it runs a Paxos round per row, so it is much slower, and rows that already existed come back in notApplied with their current values.

Delete Rows​

FieldDefaultDescription
Table-The table to delete from – required
Key-JSON object of primary key columns – required
If ExistsoffDelete only if the row exists, and report whether it did. Needs the full primary key
Consistency, TimeoutAs above

Naming only the partition key deletes the whole partition.

Execute Statement​

Run any single CQL statement with bound values: an UPDATE of a counter, an INSERT … IF NOT EXISTS, a CREATE TABLE. USE is refused, because the keyspace belongs to the connection.

A lightweight transaction (IF …) reports whether it applied; when its condition did not hold, the result carries the row as it is. A statement that returns rows returns up to 1,000 of them — run larger reads as a Query, which pages.

Execute Batch​

FieldDefaultDescription
Statements-JSON array of up to 100 {"statement": "…", "values": [...]} – required
Batch Typeloggedlogged (atomic), unlogged (faster, not atomic across partitions) or counter (counter updates)
Consistency, TimeoutAs above

Each statement must be an INSERT, UPDATE or DELETE. A value that does not fit fails the whole batch before it is sent.

A batch is not a bulk load

A batch across many partitions makes one coordinator do the work of many and slows the cluster down. Use a logged batch to keep a few writes atomic, ideally in one partition. To write many rows, use Write Rows.

Using Parameters​

Any text field accepts ((parameter)) placeholders filled from the upstream node:

FieldExample
Values["((deviceId))", "((since))"]
Key{"device_id": "((deviceId))"}
Rows((readings)), or {"device_id": "((device))", "celsius": ((value))}
Tablereadings_((plant))

A templated value is only known at run time, so the form defers its checks on fields that contain a placeholder; the same checks run on the resolved value when the function executes.

Pipeline Integration​

Each function appears as its own node under Database in the node library — see Apache Cassandra nodes for each node's output.

Common patterns:

  • Store readings — an MQTT or OPC UA trigger followed by Write Rows
  • Enrich — a Read Rows lookup before a decision node
  • Count — Execute Statement with UPDATE counters SET cycles = cycles + 1 WHERE device_id = ?
  • Keep two tables in step — a logged Execute Batch that writes a reading and updates the device's latest value together

Store-and-forward and replay safety​

OperationBuffered by store-and-forwardOn replay
Write RowsyesWrites the same rows again — same end state. With Insert Only, the first write stands
Delete RowsyesDeletes rows that are already gone — same end state
Execute Batchno— its statements may increment counters or append to lists, which would act twice
Execute Statementno— an arbitrary statement is never replayed

A write that fails part-way through a list of rows names the row that failed; the rows before it were written, and a retry writes them again with the same values.

Common Use Cases​

Time-series readings​

A table keyed PRIMARY KEY ((device_id, day), ts) keeps each device's day in one partition. An MQTT trigger writes each reading with Write Rows and a 90-day Time to Live; a dashboard pipeline reads one device's day with Read Rows.

Machine state​

A table keyed by device_id holds each machine's current state. Write Rows overwrites it on every change, and Read Rows gives any pipeline the latest record.

Counters​

A counter table counts cycles, alarms or parts per machine. Execute Statement increments it; a counter batch updates several counters together.

Troubleshooting​

SymptomCauseFix
keyspace "…" does not existThe Keyspace field names a keyspace the cluster does not haveCheck the name. A quoted, mixed-case keyspace must be written in double quotes
the cluster requires authenticationThe cluster runs PasswordAuthenticator and the connection has no usernameSet Username and Password
Provided username … and/or password are incorrectWrong credentialsCheck the role and its password
no connections were made / connection refusedNothing answers on the contact points and portCheck the addresses, the port (9042), and that the nodes' native_transport is enabled
A certificate errorTLS on, and the nodes' certificate is not signed by the CA givenProvide the right CA Certificate, or the TLS Server Name the certificate was issued for
table … does not existA typo in the Table field, or a table in another keyspaceWrite keyspace.table
has no column "…"A JSON field that is not a columnCheck the field name; the message lists the table's columns
missing partition key columnA Read Rows or Delete Rows key without the full partition keyName every partition key column
Cannot achieve consistency level …Not enough replicas alive for the levelLower the level, or bring nodes back. Transient: Retry on Fail retries it
Operation timed out / write timeoutThe replicas did not answer in timeUsually load. Transient: Retry on Fail retries it. Check the cluster's health
query must be a SELECTA write or schema statement in a QueryUse Execute Statement
ALLOW FILTERING errorA query that does not name the partition keyRestrict by the partition key, or model a table for this query