|
| 1 | +# Ecto Projections |
| 2 | + |
| 3 | +Ecto projections allow you to build read models from domain events using Ecto as the database layer. They provide automatic idempotency guarantees, transaction support, and efficient batch processing. |
| 4 | + |
| 5 | +## What are Ecto Projections? |
| 6 | + |
| 7 | +An Ecto projection is a specialized event handler that projects domain events into a relational database using Ecto. Unlike regular event handlers, Ecto projections: |
| 8 | + |
| 9 | +- **Guarantee idempotency** - Each event is projected exactly once using watermark-based tracking |
| 10 | +- **Use database transactions** - All operations execute atomically via `Ecto.Multi` |
| 11 | +- **Support batch processing** - Process multiple events in a single transaction for high throughput |
| 12 | + |
| 13 | +## Architecture |
| 14 | + |
| 15 | +### Event Handler Foundation |
| 16 | + |
| 17 | +Ecto projections are built on top of `Commanded.Event.Handler`, which means they: |
| 18 | + |
| 19 | +- Run as supervised GenServer processes |
| 20 | +- Subscribe to the event store |
| 21 | +- Receive events in order (within a single handler instance) |
| 22 | +- Support the same lifecycle callbacks (init, error handling, etc.) |
| 23 | + |
| 24 | +### Idempotency Mechanism |
| 25 | + |
| 26 | +Ecto projections use a **watermark-based idempotency** strategy: |
| 27 | + |
| 28 | +1. A `projection_versions` table tracks the last seen event number for each projector |
| 29 | +2. Before processing an event, the projector checks: `event_number > last_seen_event_number` |
| 30 | +3. If true, the event is processed and the watermark is updated |
| 31 | +4. If false, the event is skipped (already processed) |
| 32 | + |
| 33 | +This approach is: |
| 34 | +- **Simple** - Single integer comparison |
| 35 | +- **Fast** - One row per projector (not one row per event) |
| 36 | +- **Correct** - Works perfectly for sequential processing |
| 37 | + |
| 38 | +### Transaction Semantics |
| 39 | + |
| 40 | +All projection operations happen within a database transaction: |
| 41 | + |
| 42 | +```elixir |
| 43 | +Ecto.Multi.new() |
| 44 | +|> Ecto.Multi.run(:track_projection_version, fn -> update_watermark() end) |
| 45 | +|> Ecto.Multi.insert(:my_data, changeset) # Your projection logic |
| 46 | +|> Repo.transaction() |
| 47 | +``` |
| 48 | + |
| 49 | +If any step fails, the entire transaction rolls back, including the watermark update. This ensures consistency. |
| 50 | + |
| 51 | +### After-Update Callbacks |
| 52 | + |
| 53 | +The `after_update/3` and `after_update_batch/2` callbacks execute **AFTER** the transaction commits: |
| 54 | + |
| 55 | +- **Side effects only** - Use for notifications, pub/sub, external API calls |
| 56 | +- **Cannot rollback** - Database changes are already committed |
| 57 | +- **Errors propagate** - But the data is already saved |
| 58 | + |
| 59 | +This design prevents long-running side effects from blocking the transaction. |
| 60 | + |
| 61 | +## Batch Processing |
| 62 | + |
| 63 | +Batch processing allows high throughput by processing multiple events in a single database transaction: |
| 64 | + |
| 65 | +```elixir |
| 66 | +use Commanded.Projections.Ecto, |
| 67 | + batch_size: 50 # Process 50 events per transaction |
| 68 | +``` |
| 69 | + |
| 70 | +**How it works:** |
| 71 | + |
| 72 | +1. Collect up to `batch_size` events from subscription |
| 73 | +2. Start transaction |
| 74 | +3. Lock projection version row (`FOR UPDATE`) |
| 75 | +4. Filter events: keep only those with `event_number > watermark` |
| 76 | +5. Update watermark to highest event number in batch |
| 77 | +6. Execute user's projection logic for all unseen events |
| 78 | +7. Commit transaction |
| 79 | + |
| 80 | +**Benefits:** |
| 81 | + |
| 82 | +- Reduced transaction overhead (1 transaction for N events instead of N transactions) |
| 83 | +- Single fsync for the entire batch |
| 84 | +- Better throughput for high-volume event streams |
| 85 | + |
| 86 | +**Trade-offs:** |
| 87 | + |
| 88 | +- Higher latency (wait for batch to fill or timeout) |
| 89 | +- Mutually exclusive with concurrency |
| 90 | +- Requires static schema prefix (no per-event dynamic schemas) |
| 91 | + |
| 92 | +## Why Concurrency Is Not Supported |
| 93 | + |
| 94 | +> #### Concurrency Not Supported {: .error} |
| 95 | +> |
| 96 | +> Ecto projections do not support `concurrency > 1` due to the watermark-based idempotency mechanism. |
| 97 | +
|
| 98 | +**The Problem:** |
| 99 | + |
| 100 | +With concurrent workers processing events in parallel: |
| 101 | + |
| 102 | +``` |
| 103 | +Time Worker 1 Worker 2 Watermark |
| 104 | +---- -------- -------- --------- |
| 105 | +T1 Event #3 Event #5 0 |
| 106 | +T2 Updates: 3 Updates: 5 5 (Worker 2 commits first) |
| 107 | +T3 Event #4 arrives 5 |
| 108 | +T4 SKIPPED (4 < 5) ❌ 5 |
| 109 | +``` |
| 110 | + |
| 111 | +Event #4 is permanently lost because the watermark already moved past it. |
| 112 | + |
| 113 | +**Why Regular Event Handlers Can Use Concurrency:** |
| 114 | + |
| 115 | +Regular event handlers don't use watermark idempotency - they rely on the event store subscription's checkpoint. With `partition_by/2`, they guarantee per-partition ordering while allowing cross-partition concurrency. |
| 116 | + |
| 117 | +**The Solution for Ecto Projections:** |
| 118 | + |
| 119 | +Use `:batch_size` instead of `:concurrency`: |
| 120 | +- Maintains event ordering |
| 121 | +- Provides high throughput |
| 122 | +- Safe with watermark idempotency |
| 123 | + |
| 124 | +## Schema Prefixes |
| 125 | + |
| 126 | +Schema prefixes allow multi-tenant projections where each tenant's data lives in a separate PostgreSQL schema: |
| 127 | + |
| 128 | +```elixir |
| 129 | +def schema_prefix(%Event{tenant: tenant}, _metadata), do: tenant |
| 130 | +``` |
| 131 | + |
| 132 | +The `projection_versions` table will be read/written in the tenant's schema, ensuring complete isolation. |
| 133 | + |
| 134 | +> #### Batch Processing Limitation {: .warning} |
| 135 | +> |
| 136 | +> Batch projectors only support **static** schema prefixes (strings), not dynamic functions. |
| 137 | +> This is because the entire batch must use the same schema - we can't mix events from |
| 138 | +> different tenants in a single transaction. |
| 139 | +
|
| 140 | +## Design Decisions |
| 141 | + |
| 142 | +### Why Watermark Instead of Per-Event Tracking? |
| 143 | + |
| 144 | +**Watermark approach:** |
| 145 | +- Storage: 1 row per projector |
| 146 | +- Lookup: Single integer comparison |
| 147 | +- Write: 1 update per event/batch |
| 148 | + |
| 149 | +**Per-event tracking:** |
| 150 | +- Storage: 1 row per projector per event (can be millions) |
| 151 | +- Lookup: Check if event_number exists in table |
| 152 | +- Write: 1 insert per event |
| 153 | + |
| 154 | +The watermark approach is simpler and more efficient for the 99% case (sequential processing). Per-event tracking would be needed only if we wanted to support concurrent processing in the future. |
| 155 | + |
| 156 | +### Why Batch Processing Over Concurrency? |
| 157 | + |
| 158 | +Batch processing provides: |
| 159 | +- Ordering guarantees (required for correctness) |
| 160 | +- High throughput (fewer transactions) |
| 161 | +- Simpler code (single watermark update) |
| 162 | + |
| 163 | +Concurrency would require: |
| 164 | +- Per-partition watermarks |
| 165 | +- Complex coordination logic |
| 166 | +- Higher storage overhead |
| 167 | + |
| 168 | +For read models, throughput via batching is sufficient. If you need parallel processing, consider multiple projectors subscribing to different streams. |
| 169 | + |
0 commit comments