# Change streams

[Wiki](../README.md) / [Operations](README.md)

**On this page**

- [Turning streams on](#turning-streams-on)
- [Reading changes](#reading-changes)
- [Events](#events)
- [Positions and resuming](#positions-and-resuming)
- [Retention](#retention)
- [Errors](#errors)
- [Performance and safety](#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](subscriptions.md); for a live read-only
copy, use a [replica database](replicas.md).

## Turning streams on

Two switches. The instance gate, which is off by default:

```text
GDB_CHANGE_STREAMS=on
```

Then each database that should stream, as an auto-commit statement:

```cypher
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](configuration.md#settings).
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

```cypher
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/*.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.

## Related articles

[Server configuration](configuration.md) · [Commit durability and recovery](durability.md) · [Checkpointing and compaction](checkpoints.md) · [Backup, restore and transfer](backups.md)
