Add a projection
A projection hands each stored event to a handler that updates data shaped for queries. This page keeps each account's balance for a list of accounts. FCQRS follows each aggregate's versions, so the handler receives every stored event, including one whose write committed after later writes. The read side explains why that matters. Projections require a SQLite or PostgreSQL journal.
Choose where the read model lives
Where the read model lives decides where the projection keeps its progress, and what a crash means for the handler:
| Read model | Register with | Progress | After a crash or restart |
|---|---|---|---|
| In memory | Projection.single FromStart or AddProjection(handler) |
in memory | the whole journal is read again |
| In the journal's SQL database | Fcqrs.transactionalProjection or AddTransactionalProjection |
committed with each change | each event is applied once |
| Elsewhere, such as a search index | Projection.single (Named "...") or AddProjection(handler, name: "...") |
stored after the handler returns | the handler can receive an event again |
For a SQL read model in the journal's database, follow Catch up projections. Its handler writes through a transaction that FCQRS commits together with the progress.
Keep a read model in memory
The handler receives an obj because the journal holds events from every aggregate and saga. Match the
event types this projection needs and ignore the rest:
Shared setup
open FCQRS.Common
open FCQRS.FSharp
module Account =
type AccountEvent =
| Opened of owner: string
| Deposited of amount: decimal
| Withdrawn of amount: decimal
| Rejected of reason: string
open Account
FCQRSFCQRS.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_450-how-to_004-add-a-projection.md_page.AccountFcqrs_450-how-to_004-add-a-projection.md_page.Account.AccountEventOpenedowner: stringstringAn abbreviation for the CLI type . Basic Types
Depositedamount: decimaldecimalAn abbreviation for the CLI type . Basic Types
WithdrawnRejectedreason: stringopen System.Collections.Concurrent
// Each account's balance, rebuilt from the journal at every start.
let balances = ConcurrentDictionary<string, decimal>()
let handle (message: obj) =
match message with
// Sender is the ID of the account that stored the event.
| :? Event<AccountEvent> as event ->
let id = string event.Sender.Value
match event.EventDetails with
| Opened _ -> balances[id] <- 0m
| Deposited amount -> balances[id] <- balances[id] + amount
| Withdrawn amount -> balances[id] <- balances[id] - amount
| _ -> ()
| _ -> ()
SystemCollectionsConcurrentbalances: ConcurrentDictionary<string,decimal>``.ctor``: unit -> unitInitializes a new instance of the class that is empty, has the default concurrency level, has the default initial capacity, and uses the default comparer for the key type.
stringAn abbreviation for the CLI type . Basic Types
decimalAn abbreviation for the CLI type . Basic Types
handle: obj -> unitmessage: objobjAn abbreviation for the CLI type . Basic Types
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>
Fcqrs_450-how-to_004-add-a-projection.md_page.Account.AccountEventevent: 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: FCQRS.Model.Data.AggregateIdGet the value of a 'Some' option. A NullReferenceException is raised if the option is 'None'.
Sender: FCQRS.Model.Data.AggregateId optionAn optional identifier for the actor that generated the event.
EventDetails: 'EventDetailsThe specific details or payload of the event.
OpenedItem: decimalGets or sets the value associated with the specified key. The key of the value to get or set. is . The property is retrieved and does not exist in the collection. The value of the key/value pair at the specified index.
Depositedamount: decimal(+): ^T1 -> ^T2 -> ^T3Overloaded addition operator The first parameter. The second parameter. The result of the operation. 2 + 2 // Evaluates to 4 "Hello " + "World" // Evaluates to "Hello World"
Withdrawn(-): ^T1 -> ^T2 -> ^T3Overloaded subtraction operator The first parameter. The second parameter. The result of the operation. 10 - 2 // Evaluates to 8
Register it with its progress in memory:
Shared setup
let register (api: IActor) =
register: IActor -> FCQRS.Projections.IProjectionapi: IActorFCQRS.Common.IActorDefines the core functionalities and context provided by the FCQRS environment to actors. This interface provides access to essential Akka.NET services and FCQRS initialization methods.
let balanceView = Fcqrs.projection api (Projection.single FromStart handle)
balanceView: FCQRS.Projections.IProjectionFCQRS.FSharp.Fcqrsprojection: IActor -> Projection -> FCQRS.Projections.IProjectionRegister the read-model projection and return the subscription stream. Starts a projection. It follows each aggregate's and saga's own sequence numbers, so it never skips a stored event. Each pass hands the events it found to the handler in journal order, so each aggregate's events arrive in sequence order; on SQLite, events of different aggregates arrive in the order they were written, and on PostgreSQL an event that commits late arrives in a later pass. With `Named` progress, a handler can see an event again after a crash. A handler that throws terminates the process. Requires a SQLite or PostgreSQL journal; an Akka event adapter must turn each journal row into one event.
api: IActorFCQRS.FSharp.ProjectionModuleConstructors for the projection-handler shapes.
single: ProjectionProgress -> (obj -> unit) -> ProjectionSingle-event handler: just update the read model (returns unit); each aggregate event is then published to subscribers as-is. The common case when every event is worth notifying.
FromStartIn memory: the projection reads the whole journal each time it starts, then follows new events. For a read model the process keeps in memory.
handle: obj -> unitusing System.Collections.Concurrent;
using static FCQRS.Common; // Event<>
// Each account's balance, rebuilt from the journal at every start.
var balances = new ConcurrentDictionary<string, decimal>();
void Handle(object message)
{
// Sender is the ID of the account that stored the event.
if (message is not Event<AccountEvent> { Sender: { } sender } stored)
return;
var id = sender.Value.ToString();
switch (stored.EventDetails)
{
case Opened: balances[id] = 0m; break;
case Deposited deposited: balances[id] += deposited.Amount; break;
case Withdrawn withdrawn: balances[id] -= withdrawn.Amount; break;
}
}
builder.Services.AddFcqrs(connectionString, "accounts")
.AddAggregate<Account>()
.AddProjection(Handle);
The handler runs for one event at a time. It receives events in the order the journal numbered them, so one account's events arrive in version order. On PostgreSQL, an event whose write commits late can arrive after higher-numbered events of other accounts. With its progress in memory, the projection reads the whole journal each time it starts, so start-up takes longer as the journal grows. Keep the read model somewhere durable when that time matters.
Keep a read model outside the journal database
Give the projection a name to store its progress in the journal database. After a restart, it resumes after the last event it recorded:
Shared setup
balanceView
let registerNamed (api: IActor) (updateSearch: obj -> unit) =
balanceView: FCQRS.Projections.IProjectionregisterNamed: IActor -> (obj -> unit) -> FCQRS.Projections.IProjectionapi: IActorFCQRS.Common.IActorDefines the core functionalities and context provided by the FCQRS environment to actors. This interface provides access to essential Akka.NET services and FCQRS initialization methods.
updateSearch: obj -> unitobjAn abbreviation for the CLI type . Basic Types
unitThe type 'unit', which has only one value "()". This value is special and always uses the representation 'null'. Basic Types
let search =
Projection.single (Named "balance-search") updateSearch
|> Fcqrs.projection api
search: FCQRS.Projections.IProjectionFCQRS.FSharp.ProjectionModuleConstructors for the projection-handler shapes.
single: ProjectionProgress -> (obj -> unit) -> ProjectionSingle-event handler: just update the read model (returns unit); each aggregate event is then published to subscribers as-is. The common case when every event is worth notifying.
NamedStored in the journal database under this name: the projection resumes where it stopped. For a read model that outlives the process. A new name reads the whole journal once.
updateSearch: obj -> unit(|>): '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
FCQRS.FSharp.Fcqrsprojection: IActor -> Projection -> FCQRS.Projections.IProjectionRegister the read-model projection and return the subscription stream. Starts a projection. It follows each aggregate's and saga's own sequence numbers, so it never skips a stored event. Each pass hands the events it found to the handler in journal order, so each aggregate's events arrive in sequence order; on SQLite, events of different aggregates arrive in the order they were written, and on PostgreSQL an event that commits late arrives in a later pass. With `Named` progress, a handler can see an event again after a crash. A handler that throws terminates the process. Requires a SQLite or PostgreSQL journal; an Akka event adapter must turn each journal row into one event.
api: IActorShared setup
search
search: FCQRS.Projections.IProjectionbuilder.Services.AddFcqrs(connectionString, "accounts")
.AddAggregate<Account>()
.AddProjection(UpdateSearch, name: "balance-search");
FCQRS records an event as handled after the handler returns. If the process stops in between, the
handler receives that event again after the restart. Make the write idempotent: for example, store each
account's last applied Version with its data, and ignore an event whose version is not newer. A
handler that writes to several stores needs that for each of them.
The name identifies the projection's stored progress. A new name reads the whole journal once. Reusing a name resumes that projection, so do not reuse one for a different read model.
Choose which events notify callers
| F# helper | C# handler result | Subscription behaviour |
|---|---|---|
Projection.single |
void |
publish every aggregate event after handling |
Projection.filtered |
Notify |
publish or suppress the handled aggregate event |
Projection.multi |
IMessageWithCID list |
publish the exact notification list returned |
Use filtering when one command produces several events but a caller should wake only after the event that completes all required read-model updates.
Wait for the projection
Fcqrs.projection returns an IProjection; the C# host registers it for dependency injection. A
client can subscribe to a correlation id and wait until this handler has handled the matching event, or
call CatchUpAsync to wait until it has handled every event stored before the call.
Read your writes shows both. Aggregate .Send waits only for the aggregate
reply; the projection is the separate read-side confirmation.
The C# host builder supports one projection per FCQRS runtime: a second AddProjection call throws
InvalidOperationException at registration. A handler may update several read models in the same
process, and a side-by-side rebuild runs as a separate process (see
Rebuild a read model).
Handle failures visibly
Do not catch an exception and carry on. A handler that throws terminates the process, so the projection cannot stop silently while the host appears healthy. The process supervisor restarts it: a projection with its progress in memory reads the journal again, and a named one resumes after the last event it recorded. When FCQRS stops the process explains the policy.
A failure to read the journal or to store progress is handled differently: FCQRS logs the error and retries with a backoff that starts at 1 second and grows to 30 seconds, plus up to 20 percent random delay. Monitor that error log; the read model stays behind until the database is reachable again. Missing journal history, such as a deleted journal row, terminates the process, because reading again cannot bring it back.
To correct derived data, follow Rebuild a read model. Do not edit the event journal to repair a projection. Background: The read side.