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
| Field | Default | Description |
|---|---|---|
| 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)
| Field | Default | Description |
|---|---|---|
| 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 |
| Port | 9042 | CQL 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)
| Field | Default | Description |
|---|---|---|
| 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 TLS | off | Encrypt the connection. The nodes must have client_encryption_options enabled |
| Skip Certificate Verification | off | Accept 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)
| Field | Default | Description |
|---|---|---|
| Consistency | LOCAL_QUORUM | How many replicas must answer each request, unless a function chooses its own |
| Connect Timeout | 10s | Bound on connecting, including the keyspace check, and on each health check (1s–300s). Each operation has its own Timeout |
Consistency levels
| Level | Replicas that must answer |
|---|---|
ONE, TWO, THREE | That many, in any data centre |
LOCAL_ONE | One in the local data centre |
QUORUM | A majority of all replicas |
LOCAL_QUORUM | A majority in the local data centre — the usual choice |
EACH_QUORUM | A majority in every data centre (writes) |
ALL | Every replica |
ANY | Any node, even one that only stores a hint (writes) |
SERIAL, LOCAL_SERIAL | Read 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 type | Write it as |
|---|---|
text, varchar, ascii | a string |
int, bigint, smallint, tinyint, counter, varint | a whole number, or its text — every digit of a 19-digit bigint is kept |
float, double, decimal | a number, or its text |
boolean | true / false |
timestamp | RFC 3339 text (2026-09-30T10:00:00Z) or epoch milliseconds |
date | "2026-09-30" |
time | "10:15:30.5" |
uuid, timeuuid | the UUID's text |
inet | the address's text |
blob | base64, or 0x-prefixed hex |
duration | "1h30m" |
list, set | an array |
map | an object (keys converted to the map's key type) |
| user-defined type | an 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.
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.
| Field | Default | Description |
|---|---|---|
| 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"] |
| Consistency | connection's | Replicas that must answer |
| Limit | 100 | Maximum rows (1–10,000). Sent to Cassandra as the page size |
| Paging State | - | Resume where a previous run stopped |
| Timeout | 30m | Bound 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.
| Field | Default | Description |
|---|---|---|
| Table | - | readings, or plant.readings – required |
| Key | - | JSON object of primary key columns — {"device_id": "press-1"} – required |
| Columns | all | Comma-separated columns to return |
| Consistency, Limit, Paging State, Timeout | As 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.
| Field | Default | Description |
|---|---|---|
| Table | - | The table to write – required |
| Rows | - | One JSON object, or an array of objects – required |
| Insert Only | off | Write 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) |
| Consistency | connection's | Replicas that must acknowledge each write |
| Timeout | 30m | Bound 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
| Field | Default | Description |
|---|---|---|
| Table | - | The table to delete from – required |
| Key | - | JSON object of primary key columns – required |
| If Exists | off | Delete only if the row exists, and report whether it did. Needs the full primary key |
| Consistency, Timeout | As 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
| Field | Default | Description |
|---|---|---|
| Statements | - | JSON array of up to 100 {"statement": "…", "values": [...]} – required |
| Batch Type | logged | logged (atomic), unlogged (faster, not atomic across partitions) or counter (counter updates) |
| Consistency, Timeout | As 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 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:
| Field | Example |
|---|---|
| Values | ["((deviceId))", "((since))"] |
| Key | {"device_id": "((deviceId))"} |
| Rows | ((readings)), or {"device_id": "((device))", "celsius": ((value))} |
| Table | readings_((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
| Operation | Buffered by store-and-forward | On replay |
|---|---|---|
| Write Rows | yes | Writes the same rows again — same end state. With Insert Only, the first write stands |
| Delete Rows | yes | Deletes rows that are already gone — same end state |
| Execute Batch | no | — its statements may increment counters or append to lists, which would act twice |
| Execute Statement | no | — 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
| Symptom | Cause | Fix |
|---|---|---|
keyspace "…" does not exist | The Keyspace field names a keyspace the cluster does not have | Check the name. A quoted, mixed-case keyspace must be written in double quotes |
the cluster requires authentication | The cluster runs PasswordAuthenticator and the connection has no username | Set Username and Password |
Provided username … and/or password are incorrect | Wrong credentials | Check the role and its password |
no connections were made / connection refused | Nothing answers on the contact points and port | Check the addresses, the port (9042), and that the nodes' native_transport is enabled |
| A certificate error | TLS on, and the nodes' certificate is not signed by the CA given | Provide the right CA Certificate, or the TLS Server Name the certificate was issued for |
table … does not exist | A typo in the Table field, or a table in another keyspace | Write keyspace.table |
has no column "…" | A JSON field that is not a column | Check the field name; the message lists the table's columns |
missing partition key column | A Read Rows or Delete Rows key without the full partition key | Name every partition key column |
Cannot achieve consistency level … | Not enough replicas alive for the level | Lower the level, or bring nodes back. Transient: Retry on Fail retries it |
Operation timed out / write timeout | The replicas did not answer in time | Usually load. Transient: Retry on Fail retries it. Check the cluster's health |
query must be a SELECT | A write or schema statement in a Query | Use Execute Statement |
ALLOW FILTERING error | A query that does not name the partition key | Restrict by the partition key, or model a table for this query |