Subscriptions
On this page
- Creating a subscription
- Selectors: trigger queues
- Payloads: full rows and keys
- Consuming
- A consumer program
- Delivery guarantees
- Starting points and replay
- Managing subscriptions
- Where state lives
- Retention and expiry
- Errors
A subscription is a named, server-tracked position in a database's change stream. Each of its consumers receives every message, with its own position, acknowledgements and redeliveries, so several services can follow the same changes independently and resume after restarts without storing anything themselves.
Subscriptions need GDB_CHANGE_STREAMS=on and a source database with
CHANGE STREAM ON. Without selectors a subscription delivers every change.
With selectors it is a trigger queue that delivers only the label,
label-property and relationship changes it names. From version 1.4.0.
Creating a subscription
ALTER DATABASE galactus SET CHANGE STREAM ON;
CREATE SUBSCRIPTION audit ON DATABASE galactus;
CREATE SUBSCRIPTION search ON DATABASE galactus
OPTIONS {start: 'snapshot', ackTimeout: '30s', maxDeliveries: 10, maxLag: '10GiB', consumerIdleTimeout: '7d'};
| Option | Default | Meaning |
|---|---|---|
start | 'now' | Where a new consumer starts: 'now', 'earliest', 'snapshot' or a change id |
payload | 'changes' | 'changes', 'keys' or 'full': what each message's row carries |
ackTimeout | '30s' | An unacknowledged message is delivered again after this |
maxDeliveries | 10 | After this many deliveries a message is dead-lettered |
maxInFlight | 1000 | Unacknowledged messages per consumer before read returns fewer |
maxLag | unset | Expire a consumer whose unread stream exceeds this size |
consumerIdleTimeout | unset | Drop a consumer that has not read for this long |
These statements run auto-commit and need ALTER DATABASE on the source
database.
Selectors: trigger queues
CREATE SUBSCRIPTION fulfilment ON DATABASE orders
FOR (n:Order) ON CREATE, DELETE, PROPERTY status UPDATE
FOR ()-[r:PAID]->() ON CREATE, PROPERTY amount
OPTIONS {payload: 'full'};
Several FOR clauses combine (any of them matches). Each one names a node
label (n:Label) or a relationship type ()-[r:TYPE]->() and the operations
that fire it:
| Operation | Fires when |
|---|---|
CREATE | A node with the label, or a relationship of the type, is created |
DELETE | It is deleted, including by DETACH DELETE |
UPDATE | Any property, or (nodes) any label, of an existing one changes |
LABEL ADDED / LABEL REMOVED | That label is added to or removed from an existing node |
PROPERTY p | Property p is created, updated or deleted (including when the entity holding it is created or deleted) |
PROPERTY p CREATE / UPDATE / DELETE | Only that kind of change to p |
- Labels at commit: labels are evaluated as they were at commit, so
removing
:Orderstill firesFOR (n:Order) ON LABEL REMOVED. - One message per entity per transaction: a queue sends one message per
entity per transaction.
eventis the first matching change, andchangeslists every matching change of that entity in the transaction. - Before-values needed:
PROPERTY p CREATE/UPDATEand thefullpayload need the stream's before-values (the default), so they are refused on a stream withbeforeValues: false. - Matching cost: matching happens when a consumer reads, never on the commit path.
- Not yet:
WHEREconditions on selectors are not supported yet.
Payloads: full rows and keys
payload | row |
|---|---|
changes (default) | null; the event already holds the change and, for updates, the before-value |
keys | {labels, keys}: the entity's labels and its unique / node-key constraint values, so a consumer can MERGE by key |
full | {before, after, keys}: every label and property before and after the transaction (before is null for a create, after for a delete) |
Full rows are captured at commit only for entities whose label or type a
keys / full subscription names; nothing else pays. In hybrid storage this
reads the entity's unchanged properties while committing: about 3 µs per
captured entity for a 7-property node, measured on a benchmark that updates 100
captured nodes per transaction. Subscriptions without keys / full add
nothing to commits.
Consuming
CALL gdb.subscription.read('audit', {consumer: 'billing', max: 100, waitMs: 30000})
YIELD deliveryId, changeId, txId, seq, deliveryCount, metadata, event;
CALL gdb.subscription.ack('audit', $deliveryIds);
CALL gdb.subscription.nack('audit', $deliveryIds, {retryIn: '10s'});
A consumer is created by its first read. read returns up to max
messages (default 100): first any that are due for redelivery, then new ones.
With waitMs (at most 25,000, below the drivers' default 30-second timeout) it
waits for new changes when none are ready,
holding no database lock while it waits. event and metadata are exactly
what gdb.stream.read returns, so code written for the
stream works unchanged. A queue's messages also carry changes, and row
follows the payload option. Reading needs the same rights as
gdb.stream.read.
A consumer program
Any Bolt driver can consume a subscription with a short loop. This one uses the native Python driver; the pattern is the same in every language:
import os
from galactus import Driver
with Driver("bolt://127.0.0.1:7687", "gdb", os.environ["GDB_PASSWORD"],
database="galactus", timeout=30) as driver:
while True:
# Waits up to 25 s; returns as soon as a matching change commits.
result = driver.execute_query(
"CALL gdb.subscription.read('fulfilment', {consumer: 'shipper', waitMs: 25000}) "
"YIELD deliveryId, event, row")
done = []
for r in result.records:
handle(r["event"], r["row"]) # your code: make it safe to repeat
done.append(r["deliveryId"])
if done:
driver.execute_query(
"CALL gdb.subscription.ack('fulfilment', $ids)", {"ids": done})
Keep waitMs below your driver's timeout (30 seconds by default). The server
caps it at 25,000 ms.
Delivery guarantees
- At least once. A message counts as handled only when acknowledged.
Unacknowledged messages come back after
ackTimeout, or straight away after anackwith noretryIn. Make consumers idempotent:changeIdis unique per change event. - Order: messages are offered in commit order, per consumer.
- Dead letters: a message delivered
maxDeliveriestimes without an ack is set aside, counted, and recorded insystem_streams. - Acknowledging an id twice, or an id from before a reset, is a harmless no-op
that reports
acked: 0. - After a server restart, each consumer resumes from its last acknowledged position. Anything delivered but not yet acknowledged is delivered again.
Starting points and replay
| Start | Delivers |
|---|---|
'now' | Changes from now on |
'earliest' | Everything the stream still holds, and with no stream retention limit, everything since the stream was turned on |
| a change id | Everything after that change, if still retained |
'snapshot' | Every existing node and relationship as node.snapshot / rel.snapshot events, then every change since the snapshot began. This is "from zero", including data that existed before the stream |
The snapshot is read in pages while writes continue, so an entity may appear
in it with a newer state than at the snapshot's start. The changes since the
start are then replayed in order. A consumer that applies messages
idempotently (by elementId or its keys) ends up exactly in step with the
database.
Replay at any time with ALTER SUBSCRIPTION … RESET.
Managing subscriptions
SHOW SUBSCRIPTIONS;
ALTER SUBSCRIPTION audit PAUSE;
ALTER SUBSCRIPTION audit RESUME;
ALTER SUBSCRIPTION audit RESET TO 'earliest'; // every consumer
ALTER SUBSCRIPTION audit RESET CONSUMER billing TO 'snapshot'; // one consumer
ALTER SUBSCRIPTION audit DROP CONSUMER billing;
DROP SUBSCRIPTION audit [IF EXISTS];
SHOW SUBSCRIPTIONS gives one row per consumer: state, position, in-flight
count, delivered, acknowledged, redelivered and dead-lettered totals, and when
it last read. Dropping a database drops its subscriptions.
Where state lives
- Definitions are kept in
system. - Consumer positions, counters and dead letters are kept in the reserved
system_streamsdatabase, written in 100 ms batches. You can query it read-only, for exampleMATCH (c:Consumer)-[:CONSUMES]->(s:Subscription) RETURN s.name, c.name, c.acked. It cannot be written, dropped or re-engined, and it does not count towards the edition's database limit. - Database names starting
system_are reserved for the server. If a user database is already calledsystem_streams, turning onGDB_CHANGE_STREAMSfails at startup with a message naming it.
Retention and expiry
Every consumer pins its stream: retention (retention, maxBytes) never
removes changes a consumer has not yet acknowledged. With maxLag, a consumer
that falls further behind than that expires instead. It stops pinning, and its
reads fail with Expired until it is reset. Without maxLag, a consumer that
stops reading holds disk indefinitely. Use consumerIdleTimeout or
DROP CONSUMER for consumers that have gone away.
When the stream starts a new history (a restore, or the stream turned off and
on), consumers fail with ResetRequired rather than continue across the gap.
Errors
| Code | When |
|---|---|
Galactus.ClientError.Subscriptions.NotFound | No subscription (or consumer) by that name |
Galactus.ClientError.Subscriptions.AlreadyExists | CREATE with a taken name |
Galactus.ClientError.Subscriptions.Paused | read on a paused subscription |
Galactus.ClientError.Subscriptions.Expired | The consumer fell behind its maxLag, or its start is no longer retained |
Galactus.ClientError.Subscriptions.ResetRequired | The source stream started a new history |
Galactus.ClientError.Subscriptions.InvalidSelector | A selector that does not parse, a WHERE condition, or a selector or payload that needs before-values the stream does not record |
Galactus.ClientError.Streams.Disabled | The gate is off, or the source database has no stream |
Related articles
Change streams · Server configuration · Database administration