Change streams

Wiki / Operations

On this page

A change stream is an opt-in, replayable record of every committed change to a database: nodes and relationships created and deleted, labels added and removed, and properties set and removed, each with its transaction, time and user. Read it from any position to feed search indexes, caches, audit trails or other systems. Streams are off by default and cost nothing until switched on.

From version 1.4.0. For server-tracked positions, acknowledgements, replay and trigger queues, use subscriptions; for a live read-only copy, use a replica database.

Turning streams on

Two switches. The instance gate, which is off by default:

GDB_CHANGE_STREAMS=on

Then each database that should stream, as an auto-commit statement:

ALTER DATABASE galactus SET CHANGE STREAM ON;
ALTER DATABASE galactus SET CHANGE STREAM ON OPTIONS {retention: '7d', maxBytes: '10GiB', beforeValues: true};
ALTER DATABASE galactus SET CHANGE STREAM OFF;
OptionDefaultMeaning
retentionunset (keep everything)Drop stream segments whose newest change is older than this, e.g. '24h', '7d'. 'none' clears it
maxBytesunset (no limit)Drop the oldest segments while the stream is larger than this, e.g. '10GiB'. 'none' clears it
beforeValuestrueRecord each changed property's previous value

Unset options take the instance defaults GDB_STREAM_RETENTION and GDB_STREAM_MAX_BYTES, see server configuration. The system database never streams. Turning a stream on while it is on with different options is refused: turn it off first, which starts a new history. While a stream is on, SET ENGINE is refused.

The command needs the ALTER DATABASE privilege. Reading a stream needs EXECUTE ADMIN PROCEDURES plus traverse and read on the whole graph, because a stream exposes every change.

Reading changes

CALL gdb.stream.earliest();                 // id, txId: before the oldest retained change
CALL gdb.stream.current();                  // id, txId: after the newest change
CALL gdb.stream.read($from, 1000);          // id, txId, seq, lastInTx, metadata, event
CALL gdb.stream.status();                   // enabled, state, epoch, earliest, current, bytes, ...

gdb.stream.read returns up to limit events (default 1000, at most 100,000) strictly after $from, in commit order. Pass the last id you processed as the next $from. The stream is written in the background, so a change appears a moment after its commit; poll read until it returns rows.

Events

event.typeFields
node.createdelementId, labels, properties
node.deletedelementId, labels, properties as they were
node.label.added / node.label.removedelementId, label, labels after
node.property.setelementId, labels, key, before, after
node.property.removedelementId, labels, key, before
rel.created / rel.deletedelementId, relType, start, end, properties
rel.property.set / rel.property.removedelementId, relType, start, end, key, before, after
schemachange: the index or constraint change

DETACH DELETE produces an explicit rel.deleted for every relationship it removes, before the node.deleted. Each transaction's events are its net effect: setting a property twice in one transaction is one event, a change back to the original value is none, and a node or relationship created and deleted in the same transaction produces no events. metadata carries commitTimeMs, the executing user and recovered.

An elementId is reused after its node or relationship is deleted. Identify entities downstream by their key properties, and treat a delete followed by a create of the same id as two different entities.

Positions and resuming

A change id is an opaque 40-character string made of the stream's epoch, the transaction number (txId, the commit-log sequence number) and the event's position in the transaction (seq). Store the last id you processed and resume from it, after a client restart or a server restart alike.

The epoch changes when history becomes discontinuous: a restore, the stream being turned off and on, the instance gate being off while the database ran, or an unrecoverable stream fault. A cursor from an earlier epoch is rejected with EpochMismatch, never silently misread. Re-read the data and continue from gdb.stream.current().

Retention

Without retention or maxBytes, segments are kept indefinitely; plan disk for it. Segments roll at GDB_STREAM_SEGMENT_BYTES (default 64 MiB), and retention removes whole sealed segments, oldest first. A cursor older than gdb.stream.earliest() gets CursorExpired.

Errors

CodeWhen
Galactus.ClientError.Streams.DisabledThe instance gate is off, or the database has no stream
Galactus.ClientError.Streams.CursorExpired$from is older than the earliest retained change
Galactus.ClientError.Streams.EpochMismatch$from belongs to an earlier stream history
Galactus.ClientError.Streams.InvalidCursor$from is not a change id

Performance and safety

  • Off means unchanged. With the gate off, or a database not opted in, no stream thread or file exists and the commit path does no extra work.
  • On, nothing extra is written to the commit log. The commit hands the stream its context in memory: commit time and user, labels of changed entities, before-values and the state of deleted entities, all already held by the transaction's undo log. There is no extra fsync, no lock and no wait on the stream. On a write-heavy benchmark (100 property updates per transaction) commit time rose about 3% with before-values and not measurably without them.
  • A background thread per database waits until each commit is durable, then appends it to <database>/stream/*.seg with buffered writes and a sync every GDB_STREAM_SYNC_MS (default 1000).
  • After a crash, anything the stream had not synced is rebuilt on open while the commit log replays, reading each entity's state just before each change, so no committed change or before-value is lost from the stream. Only the commit time and user of those last transactions are unknown: their events have recovered: true and a null commitTimeMs.
  • A slow stream never slows writers. Past GDB_STREAM_MAX_LAG_BYTES of backlog it stops and reports failed in gdb.stream.status(), and the next open starts a new epoch. A stream fault never stops its database.
  • Commit-log compaction leaves alone any segment the stream may still need, so a database with a stream on can use more disk until the stream catches up.

Server configuration · Commit durability and recovery · Checkpointing and compaction · Backup, restore and transfer

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