Catch up projections
A month-end report can require every change already saved across the journal, including changes to
other accounts. Call CatchUpAsync after the aggregate reply, then query the transactional
projection's read model:
Shared setup
module Account =
open FCQRS.Common
// What a caller can ask an account to do.
type AccountCommand =
| Open of owner: string
| Deposit of amount: decimal
| Withdraw of amount: decimal
// What the account replies. Rejected is a reply only: it is never stored.
type AccountEvent =
| Opened of owner: string
| Deposited of amount: decimal
| Withdrawn of amount: decimal
| Rejected of reason: string
// What the account knows now, rebuilt from its events.
type AccountState = { Owner: string option; Balance: decimal }
// The state before the account's first event.
let initial = { Owner = None; Balance = 0m }
// Chooses what to do with a command, based on the current state.
let decide (command: Command<AccountCommand>) (state: AccountState) =
match command.CommandDetails, state.Owner with
| Open _, Some _ -> DeferEvent(Rejected "The account is already open")
| Open owner, None -> PersistEvent(Opened owner)
| _, None -> DeferEvent(Rejected "The account is not open")
| (Deposit amount | Withdraw amount), _ when amount <= 0m ->
DeferEvent(Rejected "The amount must be positive")
| Deposit amount, _ -> PersistEvent(Deposited amount)
| Withdraw amount, _ when amount > state.Balance ->
DeferEvent(Rejected $"Insufficient funds: {state.Balance} available")
| Withdraw amount, _ -> PersistEvent(Withdrawn amount)
// Applies one event to the state. A rejection changes nothing.
let fold (event: Event<AccountEvent>) (state: AccountState) =
match event.EventDetails with
| Opened owner -> { state with Owner = Some owner }
| Deposited amount -> { state with Balance = state.Balance + amount }
| Withdrawn amount -> { state with Balance = state.Balance - amount }
| Rejected _ -> state
open System
open FCQRS.Model.Data
open FCQRS.Common
open FCQRS.FSharp
open Account
open System.Threading
open FCQRS.Projections
let depositAndCatchUp (accounts: AggregateHandle<AccountCommand, AccountEvent>)
(projection: IProjection) (cid: CID) (alice: AggregateId)
(cancellationToken: CancellationToken) = async {
Fcqrs_450-how-to_006-catch-up-projections.md_page.AccountFCQRSFCQRS.CommonContains common types like Events and Commands Functionality for Write Side.
Fcqrs_450-how-to_006-catch-up-projections.md_page.Account.AccountCommandOpenowner: stringstringAn abbreviation for the CLI type . Basic Types
Depositamount: decimaldecimalAn abbreviation for the CLI type . Basic Types
WithdrawFcqrs_450-how-to_006-catch-up-projections.md_page.Account.AccountEventOpenedDepositedWithdrawnRejectedreason: stringFcqrs_450-how-to_006-catch-up-projections.md_page.Account.AccountStateOwner: string optionoptionThe type of optional values. When used from other CLI languages the empty option is the null value. Use the constructors Some and None to create values of this type. Use the values in the Option module to manipulate values of this type, or pattern match against the values directly. 'None' values will appear as the value null to other CLI languages. Instance methods on this type will appear as static methods to other CLI languages due to the use of null as a value representation. Options
Balance: decimalinitial: AccountStateNoneThe representation of "No value"
decide: Command<AccountCommand> -> AccountState -> EventAction<AccountEvent>command: Command<AccountCommand>FCQRS.Common.Command`1Represents a command to be processed by an aggregate actor. <typeparam name="'CommandDetails">The specific type of the command payload.</typeparam>
state: AccountStateCommandDetails: 'CommandDetailsThe specific details or payload of the command.
SomeThe representation of "Value of type 'T" The input value. An option representing the value.
DeferEventPublish and fold the event in the live actor without storing it or incrementing the persisted version. A deferred event does not start a saga; running sagas still receive it.
PersistEventPersist the event to the journal. The actor's state will be updated using the event handler *after* persistence succeeds.
(<=): 'T -> 'T -> boolStructural less-than-or-equal comparison The first parameter. The second parameter. The result of the comparison. 5 <= 1 // Evaluates to false 5 <= 5 // Evaluates to true [1; 5] <= [1; 6] // Evaluates to true
(>): 'T -> 'T -> boolStructural greater-than The first parameter. The second parameter. The result of the comparison. 5 > 1 // Evaluates to true 5 > 5 // Evaluates to false (1, "a") > (1, "z") // Evaluates to false
fold: Event<AccountEvent> -> AccountState -> AccountStateevent: Event<AccountEvent>FCQRS.Common.Event`1Represents an event generated by an aggregate actor as a result of processing a command. <typeparam name="'EventDetails">The specific type of the event payload.</typeparam>
EventDetails: 'EventDetailsThe specific details or payload of the event.
SystemModelFCQRS.Model.DataFCQRS.FSharpIdiomatic-F# functional facade for FCQRS. Gives F# consumers the same one-call ergonomics the C# host-builder (HostExtensions.fs) gives C#, but with F# idioms: records-of-functions for the definitions, typed handles for the results, an explicit wiring pipeline, and plain helpers for saga side effects. It is a *pure addition* that wraps only the existing primitives (IActor.InitializeActor / SagaBuilder.initSimple / Projections.startTracked / InitializeSagaStarter / CreateCommandSubscription / Actor.api) and changes nothing in the C# interop layer or the core. open FCQRS.FSharp let api = Fcqrs.actor config loggerFactory (Some (Fcqrs.connect DBType.Sqlite conn)) "Cluster" let documents = Fcqrs.aggregate api { Name="Document"; Initial=...; Decide=...; Fold=... } let slugs = Fcqrs.aggregate api { Name="Slug"; Initial=...; Decide=...; Fold=... } let publication = Fcqrs.saga api (publicationDef documents.Factory slugs.Factory) Fcqrs.wireSagaStarters api [ publication ] let subs = Fcqrs.projection api (Projection.single 0 updateReadModel) // (Projection.multi when you must control which notifications publish) // send a command and await the matching aggregate reply: let! ev = documents.Send (Fcqrs.newCid()) (Fcqrs.aggregateId id) cmd (fun e -> ...)
ThreadingFCQRS.ProjectionsTransactional projections with a journal-wide, durable catch-up boundary.
depositAndCatchUp: AggregateHandle<AccountCommand,AccountEvent> -> IProjection -> CID -> AggregateId -> CancellationToken -> Async<unit>accounts: AggregateHandle<AccountCommand,AccountEvent>FCQRS.FSharp.AggregateHandle`2What you get back after registering an aggregate.
projection: IProjectionFCQRS.Projections.IProjectionA running projection and its request-scoped notification subscriptions. Catch-up covers every application persistence ID in one committed journal snapshot. Akka's own cluster-sharding records (IDs starting with "/sharding/") are excluded. It does not wait for other projections, later writes, or external effects. A correlation-ID subscription (`Subscribe(cid, ...)`, which `sendAwaiting` uses) receives its notifications after the whole snapshot containing them commits, so it can wait longer than the commit of its own event. If the projection fails before then, the subscription is cancelled. Other subscriptions receive each notification when its event commits.
cid: CIDFCQRS.Model.Data.CIDCorrelationID for commands and Sagas. A CID cannot contain '~'. FCQRS joins a CID to an aggregate ID with '~' to name sagas and pub-sub topics, and reads the CID back from the text after the last '~', so a CID containing one never reaches its saga. Creating such a CID throws ArgumentException; the constructor throws instead of returning an error so that FCQRS releases built against the previous signature keep working. Deserialization does not run this check, so IsValid reports it for a CID read from a message.
alice: AggregateIdFCQRS.Model.Data.AggregateIdcancellationToken: CancellationTokenSystem.Threading.CancellationTokenPropagates notification that operations should be canceled.
async: AsyncBuilderBuilds an asynchronous workflow using computation expression syntax. let sleepExample() = async { printfn "sleeping" do! Async.Sleep 10 printfn "waking up" return 6 } sleepExample() |> Async.RunSynchronously
let! reply = accounts.Send cid alice (Deposit 100m) (fun _ -> true)
do! projection.CatchUpAsync(cancellationToken) |> Async.AwaitTask
// Inspect the reply, then query the balances of every account.
reply: Event<AccountEvent>accounts: AggregateHandle<AccountCommand,AccountEvent>Send: CID -> AggregateId -> 'Command -> ('Event -> bool) -> Async<Event<'Event>>Send a command and await the first matching aggregate event. Fails with SagaAlreadyStartedException, with nothing stored, when the command would start a saga that its correlation ID already started on this aggregate.
cid: CIDalice: AggregateIdDepositprojection: IProjectionCatchUpAsync: CancellationToken -> Tasks.TaskCaptures a fixed journal snapshot and waits for this projection to commit every event through it. Cancellation/timeout stops the caller's wait, not a transaction already being processed. Call after the aggregate persistence acknowledgement. The caller's ambient TransactionScope is suppressed; projection commits are independent.
cancellationToken: CancellationToken(|>): 'T1 -> ('T1 -> 'U) -> 'UApply a function to a value, the value being on the left, the function on the right The argument. The function. The function result. let doubleIt x = x * 2 3 |> doubleIt // Evaluates to 6
Microsoft.FSharp.Control.FSharpAsyncHolds static members for creating and manipulating asynchronous computations. See also F# Language Guide - Async Workflows. Async Programming
AwaitTask: Tasks.Task -> Async<unit>Return an asynchronous computation that will wait for the given task to complete and return its result. The task to await. If an exception occurs in the asynchronous computation then an exception is re-raised by this function. If the task is cancelled then is raised. Note that the task may be governed by a different cancellation token to the overall async computation where the AwaitTask occurs. In practice you should normally start the task with the cancellation token returned by let! ct = Async.CancellationToken, and catch any at the point where the overall async is started. Awaiting Results
var reply = await accounts(_ => true, cid, alice, new Deposit(100m));
await projection.CatchUpAsync(cancellationToken);
// Inspect the reply, then query the balances of every account.
These transactional projection APIs are available in FCQRS 6.4.0 and later.
CatchUpAsync captures one journal snapshot: the highest committed sequence number for each
persistence identity visible in one database read. A persistence identity identifies one actor's
journal history. The call succeeds after this projection has committed every event through those
captured sequence numbers. The snapshot includes aggregate events and saga journal records. It
excludes Akka's cluster-sharding records, whose persistence identities start with /sharding/.
Akka deletes their early history itself, and they are not application events.
The order matters. Await the aggregate reply before calling CatchUpAsync so the snapshot includes
that event when the reply was journaled. It also includes every other event committed before the
snapshot, across all application persistence identities in this journal. Writes arriving after the
snapshot do not extend this call's target.
To stay cheap on a large journal, most snapshots read only writes numbered after what the projection
handled LateWriteWindow earlier, 30 seconds by default. A write can take its number some time
before it commits. One that took longer than the window is handled by the next full scan, within
FullScanInterval, and a CatchUpAsync call in between can return before it. No event is skipped;
only that call's promise narrows for such a write.
The checkpoint is durable, so this wait does not require a correlation subscription before sending. It also works after a deferred reply, although a deferred reply itself adds no journal event. Check the command outcome separately: projection completion does not turn a rejected command into a successful one.
The examples use the account from the tutorial's withdraw money step. The tutorial's statement is a transactional projection too.
Register a transactional projection
Use Fcqrs.transactionalProjection in F# or AddTransactionalProjection in C#. These are separate
from the registration in Add a projection, whose handler writes outside the
transaction. The transactional runner stores one checkpoint per persistence identity under a stable
projection name.
The application supplies a SQL connection factory and a handler. FCQRS opens a transaction, passes its connection and transaction to the handler, records progress, and commits them together. Return from the handler only when its writes have completed. Use the supplied transaction for every read-model change covered by this projection.
Use the Balances table from Add a projection. The following SQLite
registration uses one database for the journal, read model, and FCQRS-owned progress tables:
Shared setup
}
open System.Data.Common
open System.Threading.Tasks
open Akka.Persistence.Query
open Dapper
open Microsoft.Data.Sqlite
open FCQRS.ProjectionStorage
let connectionString = "Data Source=app.db;"
let configuration = Microsoft.Extensions.Configuration.ConfigurationBuilder().Build()
let loggerFactory = Microsoft.Extensions.Logging.LoggerFactory.Create(fun _ -> ())
module Sqlite =
let api = Fcqrs.actor configuration loggerFactory
(Some(Fcqrs.connect FCQRS.Actor.DBType.Sqlite connectionString)) "accounts"
SystemDataCommonThreadingTasksAkkaPersistenceQueryDapperMicrosoftSqliteFCQRSFCQRS.ProjectionStorageSQL storage for journal-wide, transactional projection catch-up.
connectionString: stringconfiguration: Extensions.Configuration.IConfigurationRoot``.ctor``: unit -> unitExtensionsConfigurationBuild: unit -> Extensions.Configuration.IConfigurationRootBuilds an with keys and values from the set of providers registered in . An with keys and values from the registered providers.
loggerFactory: Extensions.Logging.ILoggerFactoryCreate: Action<Extensions.Logging.ILoggingBuilder> -> Extensions.Logging.ILoggerFactoryCreates new instance of configured using provided delegate. A delegate to configure the . The that was created.
LoggingMicrosoft.Extensions.Logging.LoggerFactoryProduces instances of classes based on the given providers.
Fcqrs_450-how-to_006-catch-up-projections.md_page.Sqliteapi: IActorFCQRS.FSharp.Fcqrsactor: Extensions.Configuration.IConfiguration -> Extensions.Logging.ILoggerFactory -> FCQRS.Actor.Connection option -> string -> IActorCreate the actor system from plain values (cluster name as a string).
SomeThe representation of "Value of type 'T" The input value. An option representing the value.
connect: FCQRS.Actor.DBType -> string -> FCQRS.Actor.ConnectionBuild a SQLite/etc. Connection from a raw connection string of any length.
SqliteSQLite using Microsoft.Data.Sqlite provider
FCQRS.ActorFCQRS.Actor.DBTypeRepresents the type of database connection
// NuGet: Dapper, Microsoft.Data.Sqlite
open System
open System.Data.Common
open System.Threading.Tasks
open Akka.Persistence.Query
open Dapper
open Microsoft.Data.Sqlite
open FCQRS.Common
open FCQRS.FSharp
open FCQRS.ProjectionStorage
open FCQRS.Projections
open Account
let openAccount =
"insert into Balances (Id, Owner, Balance) values (@Id, @Owner, 0)"
let changeBalance =
"update Balances set Balance = Balance + @Amount where Id = @Id"
let handle (connection: DbConnection) (transaction: DbTransaction)
(envelope: EventEnvelope) : Task =
task {
match envelope.Event with
// Sender is the ID of the account that stored the event.
| :? Event<AccountEvent> as stored ->
let id = string stored.Sender.Value
let add (amount: decimal) =
task {
let row = {| Id = id; Amount = amount |}
let! rows =
connection.ExecuteAsync(changeBalance, row, transaction)
if rows <> 1 then
failwith "A balance change needs an earlier Opened"
}
match stored.EventDetails with
| Opened owner ->
let row = {| Id = id; Owner = owner |}
let! _ = connection.ExecuteAsync(openAccount, row, transaction)
()
| Deposited amount -> do! add amount
| Withdrawn amount -> do! add -amount
| Rejected _ -> ()
| _ -> ()
} :> Task
let store =
SqlProjectionStore(
ProjectionSqlDialect.Sqlite,
Func<DbConnection>(fun () -> new SqliteConnection(connectionString)))
let options = TransactionalProjectionOptions("Balances", store)
let projection = Fcqrs.transactionalProjection api options handle
SystemDataCommonThreadingTasksAkkaPersistenceQueryDapperMicrosoftSqliteFCQRSFCQRS.CommonContains common types like Events and Commands Functionality for Write Side.
FCQRS.FSharpIdiomatic-F# functional facade for FCQRS. Gives F# consumers the same one-call ergonomics the C# host-builder (HostExtensions.fs) gives C#, but with F# idioms: records-of-functions for the definitions, typed handles for the results, an explicit wiring pipeline, and plain helpers for saga side effects. It is a *pure addition* that wraps only the existing primitives (IActor.InitializeActor / SagaBuilder.initSimple / Projections.startTracked / InitializeSagaStarter / CreateCommandSubscription / Actor.api) and changes nothing in the C# interop layer or the core. open FCQRS.FSharp let api = Fcqrs.actor config loggerFactory (Some (Fcqrs.connect DBType.Sqlite conn)) "Cluster" let documents = Fcqrs.aggregate api { Name="Document"; Initial=...; Decide=...; Fold=... } let slugs = Fcqrs.aggregate api { Name="Slug"; Initial=...; Decide=...; Fold=... } let publication = Fcqrs.saga api (publicationDef documents.Factory slugs.Factory) Fcqrs.wireSagaStarters api [ publication ] let subs = Fcqrs.projection api (Projection.single 0 updateReadModel) // (Projection.multi when you must control which notifications publish) // send a command and await the matching aggregate reply: let! ev = documents.Send (Fcqrs.newCid()) (Fcqrs.aggregateId id) cmd (fun e -> ...)
FCQRS.ProjectionStorageSQL storage for journal-wide, transactional projection catch-up.
FCQRS.ProjectionsTransactional projections with a journal-wide, durable catch-up boundary.
Fcqrs_450-how-to_006-catch-up-projections.md_page.AccountopenAccount: stringchangeBalance: stringhandle: DbConnection -> DbTransaction -> EventEnvelope -> Taskconnection: DbConnectionSystem.Data.Common.DbConnectionDefines the core behavior of database connections and provides a base class for database-specific connections.
transaction: DbTransactionSystem.Data.Common.DbTransactionDefines the core behavior of database transactions and provides a base class for database-specific transactions.
envelope: EventEnvelopeAkka.Persistence.Query.EventEnvelopeEvent wrapper adding meta data for the events in the result stream of query, or similar queries. The is the time the event was stored, in ticks. The value of this property represents the number of 100-nanosecond intervals that have elapsed since 12:00:00 midnight, January 1, 0001 in the Gregorian calendar (same as `DateTime.Now.Ticks`).
System.Threading.Tasks.TaskRepresents an asynchronous operation.
task: TaskBuilderBuilds a task using computation expression syntax.
Event: objFCQRS.Common.Event`1Represents an event generated by an aggregate actor as a result of processing a command. <typeparam name="'EventDetails">The specific type of the event payload.</typeparam>
Fcqrs_450-how-to_006-catch-up-projections.md_page.Account.AccountEventstored: Event<AccountEvent>id: stringstring: 'T -> stringConverts the argument to a string using ToString. For standard integer and floating point values and any type that implements IFormattable, ToString conversion uses CultureInfo.InvariantCulture. The input value. The converted string. string 'A' // evaluates to "A" string 0xff // evaluates to "255" string -10 // evaluates to "-10"
Value: AggregateIdGet the value of a 'Some' option. A NullReferenceException is raised if the option is 'None'.
Sender: AggregateId optionAn optional identifier for the actor that generated the event.
add: decimal -> Task<unit>amount: decimaldecimalAn abbreviation for the CLI type . Basic Types
row: {| Amount: decimal; Id: string |}Id: stringAmount: decimalrows: intExecuteAsync: string * obj * Data.IDbTransaction * Nullable<int> * Nullable<Data.CommandType> -> Task<int>Execute a command asynchronously using Task. The connection to query on. The SQL to execute for this query. The parameters to use for this query. The transaction to use for this query. Number of seconds before command execution timeout. Is it a stored proc or a batch? The number of rows affected.
(<>): 'T -> 'T -> boolStructural inequality The first parameter. The second parameter. The result of the comparison. 5 <> 5 // Evaluates to false 5 <> 6 // Evaluates to true [1; 2] <> [1; 2] // Evaluates to false
failwith: string -> 'TThrow a exception. The exception message. Never returns. let failingFunction() = failwith "Oh no" // Throws an exception true // Never reaches this failingFunction() // Throws a System.Exception
EventDetails: 'EventDetailsThe specific details or payload of the event.
Openedowner: stringrow: {| Id: string; Owner: string |}Owner: stringDepositedWithdrawn(~-): ^T -> ^TOverloaded unary negation. The value to negate. The result of the operation.
Rejectedstore: SqlProjectionStore``.ctor``: ProjectionSqlDialect * Func<DbConnection> -> SqlProjectionStoreUses one database for both the journal and the transactional read model.
FCQRS.ProjectionStorage.ProjectionSqlDialectDatabase SQL syntax supported by the transactional projection store.
Sqlite: ProjectionSqlDialectSQLite with a provider such as Microsoft.Data.Sqlite.
System.Func`1Encapsulates a method that has no parameters and returns a value of the type specified by the parameter. The type of the return value of the method that this delegate encapsulates. The return value of the method that this delegate encapsulates.
Microsoft.Data.Sqlite.SqliteConnectionRepresents a connection to a SQLite database. Connection Strings Async Limitations
connectionString: stringoptions: TransactionalProjectionOptions``.ctor``: string * SqlProjectionStore -> TransactionalProjectionOptionsprojection: IProjectionFCQRS.FSharp.FcqrstransactionalProjection: IActor -> TransactionalProjectionOptions -> (DbConnection -> DbTransaction -> EventEnvelope -> Task) -> IProjectionStarts a transactional projection with journal-wide CatchUpAsync support. Write the read model through the supplied connection and transaction; FCQRS commits those updates and contiguous per-persistence-ID progress together. Retain unprocessed journal history. Each pass applies events in journal order: per persistence ID in sequence order, and across persistence IDs in write order on SQLite; on PostgreSQL an event that commits late is applied in a later pass. See TransactionalProjectionOptions.
api: IActor// NuGet: Dapper, Microsoft.Data.Sqlite
using System.Data.Common;
using Akka.Persistence.Query;
using Dapper;
using Microsoft.Data.Sqlite;
using FCQRS;
using static FCQRS.Common;
using static FCQRS.ProjectionStorage;
using static FCQRS.Projections;
const string OpenAccount =
"insert into Balances (Id, Owner, Balance) values (@Id, @Owner, 0)";
const string ChangeBalance =
"update Balances set Balance = Balance + @Amount where Id = @Id";
static async Task Handle(
DbConnection connection, DbTransaction transaction, EventEnvelope envelope)
{
// Sender is the ID of the account that stored the event.
if (envelope.Event is not Event<AccountEvent> { Sender: { } sender } stored)
return;
var id = sender.Value.ToString();
async Task Add(decimal amount)
{
var row = new { Id = id, Amount = amount };
var rows =
await connection.ExecuteAsync(ChangeBalance, row, transaction);
if (rows != 1)
throw new InvalidOperationException(
"A balance change needs an earlier Opened");
}
switch (stored.EventDetails)
{
case Opened opened:
var owner = new { Id = id, opened.Owner };
await connection.ExecuteAsync(OpenAccount, owner, transaction);
break;
case Deposited deposited: await Add(deposited.Amount); break;
case Withdrawn withdrawn: await Add(-withdrawn.Amount); break;
}
}
var store = new SqlProjectionStore(
ProjectionSqlDialect.Sqlite, () => new SqliteConnection(connectionString));
var options = new TransactionalProjectionOptions("Balances", store);
builder.Services.AddFcqrs(connectionString, "accounts")
.AddAggregate<Account>()
.AddTransactionalProjection(options, Handle);
The C# registration adds IProjection to dependency injection. Inject it into the caller that waits
for catch-up. The F# facade returns the same interface. It also implements ISubscribe, so existing
correlation subscriptions remain available, with aggregate notifications published after commit.
The runner applies each read of the journal in journal order, but an event whose write commits late
is applied in a later read, so a saga's follow-up event can commit before the event that caused it.
Both carry the same correlation ID. A subscription
for that correlation ID, such as Subscribe(cid, ...) or sendAwaiting, therefore receives its
notifications after the whole snapshot commits. A waiter woken by the follow-up can read the event
that caused it. It can also wait longer than the commit of its own event, and it is cancelled if the
projection fails before the snapshot commits. A subscription without a correlation ID, including a
filter on the correlation ID, receives each notification when its event commits.
Subscriptions belong to the local worker. If several workers share the same projection name and
store, a worker can observe progress committed by another worker without publishing that other
worker's notifications. Use CatchUpAsync for the durable completion boundary across those workers.
The factory must return a new, unopened connection. For separate journal and read-model databases,
use the constructor taking journalConnectionFactory and projectionConnectionFactory. Both
databases must use the selected dialect. The journal factory must read the authoritative database;
the projection factory must open the store updated by the handler. Journal schema, table, and
persistence-ID and sequence-number column overrides must match the Akka journal configuration.
The runner validates the effective HOCON settings for both SQL journal readers and writers.
DataOptionsSetup and MultiDataOptionsSetup overrides are unsupported because they can replace
those settings. If akka.persistence.query.journal.sql.write-plugin is set, it must identify the
active write journal. Configured Akka event adapters are also unsupported by this reader.
For PostgreSQL, add the Npgsql package to the application, choose ProjectionSqlDialect.PostgreSql,
and configure the Akka journal for the same database. The handler above uses SQL accepted by both
SQLite and PostgreSQL:
Shared setup
module PostgreSql =
let handle = Sqlite.handle
Fcqrs_450-how-to_006-catch-up-projections.md_page.PostgreSqlhandle: DbConnection -> DbTransaction -> EventEnvelope -> TaskFcqrs_450-how-to_006-catch-up-projections.md_page.Sqliteopen Npgsql
let store =
SqlProjectionStore(
ProjectionSqlDialect.PostgreSql,
Func<DbConnection>(fun () -> new NpgsqlConnection(connectionString)))
let options = TransactionalProjectionOptions("Balances", store)
let connection = Fcqrs.connect FCQRS.Actor.DBType.PostgreSQL15 connectionString
let api = Fcqrs.actor configuration loggerFactory (Some connection) "accounts"
let projection = Fcqrs.transactionalProjection api options handle
Npgsqlstore: SqlProjectionStore``.ctor``: ProjectionSqlDialect * Func<DbConnection> -> SqlProjectionStoreUses one database for both the journal and the transactional read model.
FCQRS.ProjectionStorage.ProjectionSqlDialectDatabase SQL syntax supported by the transactional projection store.
PostgreSql: ProjectionSqlDialectPostgreSQL with a provider such as Npgsql.
System.Func`1Encapsulates a method that has no parameters and returns a value of the type specified by the parameter. The type of the return value of the method that this delegate encapsulates. The return value of the method that this delegate encapsulates.
System.Data.Common.DbConnectionDefines the core behavior of database connections and provides a base class for database-specific connections.
Npgsql.NpgsqlConnectionThis class represents a connection to a PostgreSQL server.
connectionString: stringoptions: TransactionalProjectionOptions``.ctor``: string * SqlProjectionStore -> TransactionalProjectionOptionsconnection: FCQRS.Actor.ConnectionFCQRS.FSharp.Fcqrsconnect: FCQRS.Actor.DBType -> string -> FCQRS.Actor.ConnectionBuild a SQLite/etc. Connection from a raw connection string of any length.
FCQRSPostgreSQL15PostgreSQL 15+
FCQRS.ActorFCQRS.Actor.DBTypeRepresents the type of database connection
api: IActoractor: Extensions.Configuration.IConfiguration -> Extensions.Logging.ILoggerFactory -> FCQRS.Actor.Connection option -> string -> IActorCreate the actor system from plain values (cluster name as a string).
configuration: Extensions.Configuration.IConfigurationRootloggerFactory: Extensions.Logging.ILoggerFactorySomeThe representation of "Value of type 'T" The input value. An option representing the value.
projection: IProjectiontransactionalProjection: IActor -> TransactionalProjectionOptions -> (DbConnection -> DbTransaction -> EventEnvelope -> Task) -> IProjectionStarts a transactional projection with journal-wide CatchUpAsync support. Write the read model through the supplied connection and transaction; FCQRS commits those updates and contiguous per-persistence-ID progress together. Retain unprocessed journal history. Each pass applies events in journal order: per persistence ID in sequence order, and across persistence IDs in write order on SQLite; on PostgreSQL an event that commits late is applied in a later pass. See TransactionalProjectionOptions.
handle: DbConnection -> DbTransaction -> EventEnvelope -> Taskvar store = new SqlProjectionStore(
ProjectionSqlDialect.PostgreSql,
() => new Npgsql.NpgsqlConnection(connectionString));
var options = new TransactionalProjectionOptions("Balances", store);
builder.Services.AddFcqrs(
connectionString, "accounts", FCQRS.Actor.DBType.PostgreSQL15)
.AddAggregate<Account>()
.AddTransactionalProjection(options, Handle);
The snapshot and completion contract uses per-identity checkpoints on both databases.
FCQRS owns the transaction and progress tables. Keep the handler's connection and transaction inside the handler, and let FCQRS commit or roll back. Dispose the F# projection handle when its runtime scope ends; the C# host manages the registered projection's lifetime.
Understand completion and ordering
Transactional projections handle events in sequence within each persistence identity. They do not promise a global processing order across identities. A handler combining facts from different aggregates must tolerate their arrival order or explicitly coordinate those dependencies.
The runner requires a complete history through each captured target. It does not treat the largest observed sequence number as proof that missing earlier events were handled. Preserve journal events needed by this projection, and begin with a new read model when starting a new checkpoint history.
The completion boundary covers one projection and the database transaction used by its handler. It does not wait for another projection, an external HTTP call, or work started without awaiting it. A query can observe later committed changes too; completion does not freeze the read model at the captured snapshot.
Catch-up suppresses an ambient TransactionScope so the snapshot sees freshly committed journal
data. Projection commits belong to FCQRS transactions independently of the caller's transaction;
rolling back the caller's scope does not roll back projection work.
For an account-closing workflow, this wait can establish that the selected projection processed every event committed before the closing event and the subsequent snapshot. Rejecting later commands to the closed account remains an aggregate or application rule. The wait does not drain commands that were sent earlier but have not yet been persisted, and it does not stop future writes.
Configure discovery and waiting
TransactionalProjectionOptions configures the runner:
| Property | Default | Meaning |
|---|---|---|
PollInterval |
1 second | Background delay before discovering new journal heads after the previous batch finishes |
BatchSize |
500 | Maximum events fetched for one persistence identity per query |
CatchUpTimeout |
30 seconds | Time allowed for the entire catch-up call, including snapshot capture |
LateWriteWindow |
30 seconds | How long after taking its journal number a write may commit and still be found by a query of recent writes |
FullScanInterval |
5 minutes | How often a query reads every history's head, which also finds a write that committed after the window |
Background discovery follows new events as the application runs. Each CatchUpAsync call captures
its own fixed target once; it does not repeatedly replace the target with a newer journal tail.
Most discovery queries read only the journal rows numbered after what the projection handled a
LateWriteWindow earlier, so their cost follows recent writes. A full scan, which groups the whole
journal by persistence_id, runs when the projection starts, once more a window later, and every
FullScanInterval. It costs more as the journal grows. An aggregate on this node that stores an
event starts the next query at once.
The catch-up timeout belongs to TransactionalProjectionOptions, separately from
akka.fcqrs.command-timeout used by Read your writes.
Recover without losing progress
Read-model updates and their checkpoint share one transaction. If the handler fails before commit, neither becomes durable. After restart, use the same projection name and store to resume from the committed checkpoints. The handler can run again for an event whose transaction did not commit; external side effects therefore need their own retry and idempotency policy.
An error while processing journal events stops the runner and faults IProjection.Completion.
Catch-up calls then report the error. Observe this task in the application's worker or health
monitoring, correct the cause, and restart the projection. Missing sequence numbers stop the runner
instead of silently advancing its checkpoint.
Cancellation or TimeoutException ends the caller's wait. It does not undo the aggregate command or read-model
transactions that have already committed. Treat the result as an incomplete confirmation, then
retry catch-up or report that the read model has not yet been confirmed current.
The repository's facade tests cover catch-up and recovery:
dotnet run --project test/Facade.Tests/Facade.Tests.fsproj
Set FCQRS_TEST_POSTGRES to a PostgreSQL test-server connection string to include the PostgreSQL
integration cases. The test account needs permission to create and drop the tests' isolated
databases. CI runs these cases against its PostgreSQL service. Test your domain
covers tests for the application's own decisions and replay rules.
For rebuilding query data, continue with Rebuild a read model. For a request that needs one matching notification, use Read your writes.