Capstone 02 — Stateful Fraud / Risk Stream Processing (runnable Scala)
Phase: 16 — Capstone Systems | Difficulty: ⭐⭐⭐⭐⭐ | Time: 1–2 weeks
The JD's Problem 2, built as a Flink-style keyed, event-time stateful operator in Scala (P04/P05/P11) — verified with
sbt test. Per-account state, a hot-swappable broadcast ruleset, three rule kinds, idempotent (deduped) alerts, and savepoint-style checkpoint/restore for stateful redeploys.
What it implements
| Capability | In the code | Phase |
|---|---|---|
| Keyed state per account | state: Map[account, AccountState] (recent txns + last country) | P04 |
| Broadcast rules (hot-swap) | updateRules(...) swaps the active ruleset | P04 |
| Amount rule | single txn over a threshold | — |
| Velocity rule | N txns within a time window (event-time) | P04 |
| Impossible travel | country switch faster than physically possible | P05 enrichment |
| Idempotent alerts | dedup by (account, kind, eventTime) → exactly-once effect | P01 |
| Savepoint redeploy | snapshotState / restore | P04 |
| Replay validation | re-run the stream → identical alert set | P02/P13 |
Run
sbt test
Expected: 8 specs, all green — each rule fires correctly, alerts dedup, a broadcast rule update changes behavior, and the headline proofs:
- savepoint/restore: processing
[1,2,3]then snapshot→restore→process[4,5]yields the same alert set as processing[1..5]in one detector (no lost/duplicated alerts across a stateful redeploy). - replay: re-running the same stream produces the identical alert set.
Success criteria
sbt test→ all 8 specs pass.- You can explain why alerts must be deduped (at-least-once delivery + idempotent effect = P01) and how a savepoint preserves the dedup set across a redeploy.
- You can explain why velocity/impossible-travel are inherently event-time, not arrival-time.
How to extend toward production
- Add watermarks + allowed lateness (Phase 04
lab-02-scala-windows) so late txns are handled, not silently misordered. - Enrich from the Cassandra serving layer (Capstone 03) for account risk profiles (the as-of join, P05/P11), via async lookups.
- Emit alerts to the lakehouse (Phase 09 lab) for audit + replay, and to a low-latency sink for action.
- Wrap the detector in an FS2 stream (P12) reading a real Kafka topic; checkpoint state to S3 (the real savepoint).