Перейти к содержимому
Русский

Persistence query

Это содержимое пока не доступно на вашем языке.

The PersistenceQuery interface is the read-side API for the journal. It’s how projections consume events; it’s also the interface for any one-shot read of historical events.

Two flavors:

FlavorReturnsUse
current*Promise<Event[]>Batch read of events at call time. Resolves once.
events*AsyncIterable<Event>Live stream. Polls until the consumer stops.
QueryPairs with
InMemoryQueryInMemoryJournal — and any journal without its own query
SqliteQuerySqliteJournal
PostgresQueryPostgresJournal
MariaDbQueryMariaDbJournal (MySQL too)
CassandraQueryCassandraJournal
MongoQueryMongoJournal
RelationalQueryany RelationalJournal — the dialect-neutral base the two SQL rows above are named subclasses of

The interface is the same across all of them — construct the query for your journal backend and pass it the journal instance, not a connection string:

import { PostgresJournal, PostgresJournalOptions, PostgresQuery } from 'actor-ts/persistence';
// PostgresQuery reads the same journal your PersistentActors write to.
const journalOptions = PostgresJournalOptions.create().withUrl(process.env.DATABASE_URL!);
const journal = new PostgresJournal(journalOptions);
const query = new PostgresQuery(journal);

InMemoryQuery works against any Journal — it is the reference implementation, not the in-memory-only one — so a backend missing from the table above is not blocked from running projections. What you lose is the index, and only on the by-tag methods:

  • currentEventsByPersistenceId and the id methods delegate straight to the journal, so they cost the same either way.
  • currentEventsByTag has no index to push the filter into. It lists every persistence id, replays each one from sequence 1, and filters in process — the whole journal, per call, with a by-tag projection repeating that every pollIntervalMs.

That is correct at any volume and fast only at small ones. A dedicated query class exists exactly where the backend keeps a tag index worth reading: the SQL backends maintain an events_tags join table, Cassandra an events_by_tag side table, MongoDB an index on the tags array.

const events = await query.currentEventsByPersistenceId<AccountEvent>(
'account-42',
/* fromSeq */ 1,
/* toSeq? */ 100,
);

One-shot read of every event for one entity, between sequence numbers fromSeq and toSeq (inclusive).

Use for:

  • Audit replays — get every event for one user.
  • One-off backfills — process events from one entity into a separate system.
for await (const event of query.eventsByPersistenceId<AccountEvent>('account-42', 1)) {
console.log(event.event.kind, event.sequenceNr);
if (someCondition) break; // stops the polling
}

Live stream. Yields past events first (chronological by sequenceNr), then new events as they’re appended. The stream never completes on its own — break out of the loop to stop.

Use for:

  • Live per-entity views — a real-time feed of one user’s activity.
  • Single-entity projections built ad-hoc rather than via ProjectionActor.byPersistenceId.
const events = await query.currentEventsByTag<Event>(
'account',
offsetStart,
);

One-shot read of every event matching a tag (or tag filter), starting from an offset.

Returns [Event, Offset][] pairs — the offset is what you’d pass to a subsequent call to read the next batch.

Use for:

  • Reports — sum all account events for the current month.
  • Batch ETL — periodic export of all tagged events.
for await (const [event, offset] of query.eventsByTag<Event>('account', offsetStart)) {
await handle(event);
await offsetStore.save('my-projection', offset);
}

Live stream of tag-matched events. This is what ProjectionActor.byTag uses internally — the offset lets the projection persist its progress and resume from there.

Three methods answer “which entities are in this journal?”, and the difference between them is what they cost and when they stop.

MethodReturnsStops
currentPersistenceIds()Promise<string[]>Immediately — the whole list, one array.
currentPersistenceIdsPaginated()AsyncIterable<string>When the journal is exhausted.
allPersistenceIds()AsyncIterable<string>Never — new ids keep arriving.
const ids = await query.currentPersistenceIds();

The whole list in one array. Convenient, and the right choice whenever the id count is small — a few hundred entities cost nothing to materialize. It is, though, the only method here with no bound on what it allocates: a million entities is a million strings.

for await (const persistenceId of query.currentPersistenceIdsPaginated({ pageSize: 500 })) {
await backfill(persistenceId);
}

The same data, fetched a page at a time, so peak memory is one page rather than the whole set. Completes once the journal is exhausted.

The cursor is a persistence id, not an opaque token — resume a partial walk by passing the last id you handled:

for await (const persistenceId of query.currentPersistenceIdsPaginated({
afterPersistenceId: lastHandled,
})) {
// ... picks up strictly after `lastHandled`
}
OptionDefaultWhat
pageSize256Ids fetched per round-trip.
afterPersistenceId—Resume after this id, exclusive.

Ids come back in the backend’s ascending order — byte-wise on SQLite and Cassandra, collation-aware on Postgres. Which one it is never matters to a consumer, because the cursor is compared in the same order the page is sorted in.

for await (const persistenceId of query.allPersistenceIds()) {
startProjectionFor(persistenceId); // runs forever
}

Live stream: every id the journal has ever seen, then each new one as it first appears. This is the fan-out primitive — start a per-entity consumer as its entity shows up, instead of polling currentPersistenceIds and diffing the result yourself.

Two properties worth knowing before you build on it:

  • Once per stream, not once per journal. A fresh subscription re-emits every id. A consumer that must not act twice across a restart needs its own checkpoint — the same at-least-once posture the event queries have.
  • The stream remembers what it emitted. That is an in-process set of ids, so a stream over a journal with a million entities holds a million strings for as long as it runs. A cheaper high-water mark would be wrong: ids are not created in sorted order, so an entity whose id sorts below the mark would never be emitted at all. For a one-shot enumeration that needs no such set, use currentPersistenceIdsPaginated.
type Offset = {
readonly timestamp: number;
readonly persistenceId: string;
readonly sequenceNr: number;
};

Compound: (timestamp, persistenceId, sequenceNr). This is the order events stream out.

  • timestamp — the journal’s write time for the event.
  • persistenceId — tiebreaker for events with the same timestamp.
  • sequenceNr — final tiebreaker for events in the same pid.

The compound order means events from different entities interleave by timestamp, while events from the same entity stay in their original sequence.

offsetStart is the sentinel “from the beginning.” Persist the offset returned alongside each event; resume from the persisted offset on restart.

import type { TagFilter } from 'actor-ts/persistence';
// All events that have at least one of these tags:
query.eventsByTag<E>({ any: ['account', 'audit'] }, offset);
// Events that have ALL these tags:
query.eventsByTag<E>({ all: ['account', 'production'] }, offset);
// AND across two groups:
query.eventsByTag<E>({
all: ['account'],
any: ['debit', 'credit'],
}, offset);

For simple cases, a single tag string works:

query.eventsByTag<E>('account', offset); // equivalent to { all: ['account'] }
const stream = query.eventsByTag<E>('account', offset, {
pollIntervalMs: 500, // default 1000
});

Live streams poll the journal. Configurable via LiveQueryOptions:

OptionDefaultWhat
pollIntervalMs1000How often to check for new events.

For sub-poll-interval delivery, see Push-based query.

When to use it directly vs via ProjectionActor

Section titled “When to use it directly vs via ProjectionActor”
// Manual iteration — full control:
for await (const [event, offset] of query.eventsByTag(...)) {
// ... custom error handling, custom dispatching ...
}
// ProjectionActor — supervised, restartable, offset-managed:
ProjectionActor.byTag(system, ByTagProjectionOptions.create() /* .withName(...).withTag(...).withQuery(...).withHandle(...) */);

ProjectionActor wraps the query API with:

  • Offset persistence (you don’t manage it).
  • Supervised restart on handler errors.
  • Coordinated shutdown awareness.
  • Configurable retry / backoff.

For production projections, use ProjectionActor. For one-off scripts and ad-hoc admin queries, use the query API directly.

  • current* reads scan the journal once, return a materialized array. For large pids / many tag-matches, this can be slow + memory-heavy. Use the streaming events* variants for unbounded reads, and currentPersistenceIdsPaginated instead of currentPersistenceIds.
  • Id pagination is pushed into the backend where the store has a sorted key over ids: an ORDER BY … LIMIT on SQLite and the SQL backends, a clustering-column range on the Cassandra all_persistence_ids table. MongoDB and DynamoDB have no such index — distinct has no cursor and a DynamoDB Scan has no order across partition keys — so there the page is cut in-process and pagination buys ergonomics, not I/O.
  • Live polling is light at idle — empty queries return quickly. Under heavy write load, polling can pile up if the read side is slower than the write side. Parallelize the handler if you see lag.
  • Tag indexes matter — events_tags index is what makes by-tag queries fast. The SQLite journal creates this automatically.
  • A compacted event leaves the tag index too. delete is specified against the read side as well as read: once a prefix is compacted, currentEventsByTag stops returning it. Backends whose index sits on the event record itself (MongoDB’s multikey tags) get that for free; the ones with a separate structure — SQLite’s and the SQL backends’ tags table, Cassandra’s events_by_tag — delete from it as part of the same call, and before the events, so a crash mid-delete cannot strand rows whose keys are no longer reconstructible. Worth knowing when you write your own journal: read and highestSeq will look perfectly correct while a forgotten index keeps serving deleted events.

The PersistenceQuery API reference covers the full surface.