Change streams
On this page
- Turning streams on
- Reading changes
- Events
- Positions and resuming
- Retention
- Errors
- Performance and safety
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;
| Option | Default | Meaning |
|---|---|---|
retention | unset (keep everything) | Drop stream segments whose newest change is older than this, e.g. '24h', '7d'. 'none' clears it |
maxBytes | unset (no limit) | Drop the oldest segments while the stream is larger than this, e.g. '10GiB'. 'none' clears it |
beforeValues | true | Record 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.type | Fields |
|---|---|
node.created | elementId, labels, properties |
node.deleted | elementId, labels, properties as they were |
node.label.added / node.label.removed | elementId, label, labels after |
node.property.set | elementId, labels, key, before, after |
node.property.removed | elementId, labels, key, before |
rel.created / rel.deleted | elementId, relType, start, end, properties |
rel.property.set / rel.property.removed | elementId, relType, start, end, key, before, after |
schema | change: 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
| Code | When |
|---|---|
Galactus.ClientError.Streams.Disabled | The 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/*.segwith buffered writes and a sync everyGDB_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: trueand a nullcommitTimeMs. - A slow stream never slows writers. Past
GDB_STREAM_MAX_LAG_BYTESof backlog it stops and reportsfailedingdb.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.
Related articles
Server configuration · Commit durability and recovery · Checkpointing and compaction · Backup, restore and transfer