Asynchronous Projections
The async daemon is a background service that processes events and applies projections asynchronously, providing eventually consistent read models.
How It Works
- The High Water Mark Detector monitors
pc_eventsfor new events using SQLLEAD()window functions to detect gaps - The Event Loader fetches batches of events for processing
- Each projection processes its batch and updates its read model
- Progress is tracked in
pc_event_progressionvia atomicMERGEstatements
Enabling the Async Daemon
Register projections with async lifecycle:
var store = DocumentStore.For(opts =>
{
opts.Connection("...");
opts.Projections.Snapshot<OrderSummary>(SnapshotLifecycle.Async);
opts.Projections.Add<DashboardProjection>(ProjectionLifecycle.Async);
});When wired up through AddPolecat(), opt the daemon into the host's lifetime explicitly with AddAsyncDaemon(DaemonMode):
builder.Services.AddPolecat(opts =>
{
opts.Connection("...");
opts.Projections.Snapshot<OrderSummary>(SnapshotLifecycle.Async);
opts.Projections.Add<DashboardProjection>(ProjectionLifecycle.Async);
})
.AddAsyncDaemon(DaemonMode.Solo) // start the daemon as IHostedService
.ApplyAllDatabaseChangesOnStartup(); // run schema migration at bootTIP
Async projections do not run unless you call AddAsyncDaemon(...). Use DaemonMode.Solo for single-node deployments and DaemonMode.HotCold for multi-node deployments where only one host should own each projection shard.
Daemon Settings
Configure daemon behavior:
opts.DaemonSettings.StaleSequenceThreshold = 1000;Polling
Unlike Marten's PostgreSQL LISTEN/NOTIFY, Polecat uses polling to detect new events:
// The daemon polls for new events at a configurable interval
// Default: 500msGraceful Shutdown and the Drain Timeout
When a projection or subscription shard is stopped, the daemon does not simply cancel it. It first tries to drain the agent: let the in-flight page of events finish being applied, then flush the shard's progression row so the next start picks up exactly where this one left off. StopAndDrainTimeout bounds how long the daemon waits for that drain on a single shard:
// The default is 5 seconds
opts.Projections.StopAndDrainTimeout = TimeSpan.FromSeconds(30);The bound applies to every stop path: stopping one agent, stopping all agents (the SIGTERM/host shutdown path), and the internal stop-if-already-running replacement that happens when an agent is reassigned.
Why you would raise it. If the drain is cut off before the progression flush lands, the shard restarts against a stale progression row and throws ProgressionProgressOutOfOrderException on its next start. Raise the timeout when in-flight batches legitimately take longer than five seconds — a large BatchSize, expensive projection code, heavy rebuild load, or a slow or contended SQL Server. This is most visible shutting down a host with a large agent universe: a database-per-tenant deployment with thousands of (projection × tenant) shards all draining inside a Kubernetes termination grace window.
TIP
A per-shard bound is only useful if the process lives long enough to spend it. Match a raised StopAndDrainTimeout with the host's own HostOptions.ShutdownTimeout and, on Kubernetes, the pod's terminationGracePeriodSeconds.
Why you would lower it. A deployment that would rather cut a wedged shard loose quickly and take the progression replay hit — to keep node failover and reassignment latency low, for instance — can set it below the default.
Opting out. Timeout.InfiniteTimeSpan, or any non-positive value, removes the separate bound so the drain is limited only by the daemon's own cancellation. Be aware that this means a genuinely wedged shard can hold up shutdown indefinitely.
Waiting for Non-Stale Data
CatchUpAsync
Wait for all projections to catch up to the current high water mark:
await store.WaitForNonStaleProjectionDataAsync(TimeSpan.FromSeconds(30));Per-Query
Wait for projections before a specific query:
var orders = await session.Query<OrderSummary>()
.QueryForNonStaleData()
.Where(x => x.Status == "Active")
.ToListAsync();Event Progression
Track daemon progress:
// The pc_event_progression table stores:
// - name: Projection/subscription name
// - last_seq_id: Last processed sequence ID
// - last_updated: When last updatedHigh Water Mark Detection
The high water mark detector uses SQL Server's LEAD() window function to detect sequence gaps in the event log. This prevents the daemon from processing events out of order when concurrent writers create gaps.
Listening for Daemon Commits
Sometimes you need to run a side effect after an aggregate has been durably updated by the async daemon — the canonical case is invalidating a cache key. Doing this before the commit would open a window where a concurrent read could repopulate the cache with stale state. Subscriptions, composite projections, and projection side effects all run before the batch is committed, so they can't close that window.
Register an IChangeListener on Projections.AsyncListeners to hook the commit boundary of each daemon projection batch:
public class CacheFlushingListener : IChangeListener
{
// Runs AFTER the batch is committed → "at most once". Ideal for cache invalidation.
public async Task AfterCommitAsync(IDocumentSession session, IChangeSet commit, CancellationToken token)
{
foreach (var party in commit.Updated.OfType<QuestParty>())
{
await _cache.RemoveAsync($"quest-party:{party.Id}", token);
}
}
// Runs BEFORE the batch is committed → "at least once".
public Task BeforeCommitAsync(IDocumentSession session, IChangeSet commit, CancellationToken token)
=> Task.CompletedTask;
}var store = DocumentStore.For(opts =>
{
opts.ConnectionString = connectionString;
opts.Projections.Snapshot<QuestParty>(SnapshotLifecycle.Async);
// Fires only within the async daemon, once per committed projection batch
opts.Projections.AsyncListeners.Add(new CacheFlushingListener());
});The commit parameter is an IChangeSet describing the projected documents written in that batch (Inserted / Updated / Deleted).
Delivery semantics mirror Marten:
AfterCommitAsyncruns once, after the transaction commits — at most once. A faulting after-commit listener is swallowed so the batch is not reprocessed (which would re-fire the side effect); the data is already durable.BeforeCommitAsyncruns before the commit — at least once. A throw here aborts the batch before anything is committed.
TIP
Async listeners are suppressed during projection rebuilds. A full replay re-applies every event, so firing post-commit side effects for each historical batch is almost never what you want.
Blue/Green Deployments with Projection Versioning
Every projection carries a Version (default 1). The version is baked into the shard's progression identity, so different versions of the same projection track their progress independently:
public class TripProjection : SingleStreamProjection<Trip, Guid>
{
public TripProjection()
{
// Bump the version when you change the projection's logic or shape
Version = 3;
}
}| Projection | Version | Progression identity (pc_event_progression.name) |
|---|---|---|
Trip | 1 | Trips:All |
Trip | 2 | Trips:V2:All |
Trip | 3 | Trips:V3:All |
Because each version has its own progression row, you can stand up a new version alongside the old one and let it rebuild from zero while the old version keeps serving reads — a blue/green deployment.
Gating Side Effects Behind the Prior Version
The catch with a blue/green rebuild is side effects. If your projection publishes messages or emits other side effects from RaiseSideEffects(), a new version replaying the full event history from zero would re-fire every side effect for events the previous version already processed — duplicate emails, duplicate messages, and so on.
Opt into the side-effect gate to prevent that:
public class TripProjection : SingleStreamProjection<Trip, Guid>
{
public TripProjection()
{
Version = 3;
// Suppress side effects for events the prior version already processed
Options.GateSideEffectsBehindPriorVersion = true;
}
}When the daemon starts a gated shard whose own progression is behind the highest prior version's persisted mark N, it:
- Replays the new version from its current position up to
Nin Rebuild mode, with side effects suppressed (the projected documents are still built — only the side effects are gated). - Hands off to Continuous execution from
N, where side effects fire normally.
The net effect: side effects fire exactly once, only for events past N — the events the previous version never saw.
Semantics and Edge Cases
- Trigger is "own progress
< prior mark". The gate keys off the new version's own persisted progression, so an interrupted warm-up resumes suppressed rather than re-emitting. Restarting the daemon after a crash mid-warm-up picks up from the recorded position with side effects still gated. - Failed warm-up pauses the shard. If the suppressed replay throws (e.g. a poison event under a pausing error policy), the shard is left paused with the exception attached, no continuous execution starts, and no side effects fire. Restarting the daemon resumes the warm-up from its persisted progress.
Version == 1and the flag-off case are inert. The gate is skipped entirely unlessVersion > 1andGateSideEffectsBehindPriorVersionis set — a v1 projection behaves exactly as it always has.SubscribeFromPresentis incompatible. A shard that subscribes from "present" ignores persisted progression, so the gate cannot reason about a prior mark. The gate is skipped and a warning is logged; use versioning-with-gate or subscribe-from-present, not both.- Overlap window. The prior mark
Nis snapshotted when the new version starts. If the old version is still running and advances pastNwhile the new version warms up, the events in(N, old_final]can be processed by both versions — an accepted duplicate window. Stop the old version (or accept the small overlap) when you need exactly-once across the cutover.
TIP
The gate suppresses side effects, not the projection write itself — the new version's documents are fully rebuilt over the entire history. Only RaiseSideEffects() output (published messages, emitted events, etc.) is held back during the warm-up.
Error Handling
The daemon uses Polly resilience pipelines for error handling. See Resiliency Policies for configuration.
Architecture
pc_events
│ (Polling)
▼
High Water Mark Detector
│ (Sequence Range)
▼
Event Loader
├──► Projection A ──► pc_doc_summary
├──► Projection B ──► pc_doc_dashboard
└──► Subscription C ──► External System
│
▼
pc_event_progression (tracks progress for all)
JasperFx provides formal support for Polecat and other Critter Stack libraries. Please check our