Subscriptions

Wiki / Operations

On this page

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'};
OptionDefaultMeaning
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
maxDeliveries10After this many deliveries a message is dead-lettered
maxInFlight1000Unacknowledged messages per consumer before read returns fewer
maxLagunsetExpire a consumer whose unread stream exceeds this size
consumerIdleTimeoutunsetDrop 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:

OperationFires when
CREATEA node with the label, or a relationship of the type, is created
DELETEIt is deleted, including by DETACH DELETE
UPDATEAny property, or (nodes) any label, of an existing one changes
LABEL ADDED / LABEL REMOVEDThat label is added to or removed from an existing node
PROPERTY pProperty p is created, updated or deleted (including when the entity holding it is created or deleted)
PROPERTY p CREATE / UPDATE / DELETEOnly that kind of change to p
  • Labels at commit: labels are evaluated as they were at commit, so removing :Order still fires FOR (n:Order) ON LABEL REMOVED.
  • One message per entity per transaction: a queue sends one message per entity per transaction. event is the first matching change, and changes lists every matching change of that entity in the transaction.
  • Before-values needed: PROPERTY p CREATE / UPDATE and the full payload need the stream's before-values (the default), so they are refused on a stream with beforeValues: false.
  • Matching cost: matching happens when a consumer reads, never on the commit path.
  • Not yet: WHERE conditions on selectors are not supported yet.

Payloads: full rows and keys

payloadrow
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 a nack with no retryIn. Make consumers idempotent: changeId is unique per change event.
  • Order: messages are offered in commit order, per consumer.
  • Dead letters: a message delivered maxDeliveries times without an ack is set aside, counted, and recorded in system_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

StartDelivers
'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 idEverything 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_streams database, written in 100 ms batches. You can query it read-only, for example MATCH (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 called system_streams, turning on GDB_CHANGE_STREAMS fails 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

CodeWhen
Galactus.ClientError.Subscriptions.NotFoundNo subscription (or consumer) by that name
Galactus.ClientError.Subscriptions.AlreadyExistsCREATE with a taken name
Galactus.ClientError.Subscriptions.Pausedread on a paused subscription
Galactus.ClientError.Subscriptions.ExpiredThe consumer fell behind its maxLag, or its start is no longer retained
Galactus.ClientError.Subscriptions.ResetRequiredThe source stream started a new history
Galactus.ClientError.Subscriptions.InvalidSelectorA 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.DisabledThe gate is off, or the source database has no stream

Change streams · Server configuration · Database administration

Planning a deployment? Review compatibility and licence setup for your instance.