Skip to content
Reliable Data Engineering
Practice problem hard streamingfeature-storeflinklow-latencyml
Practise with timer, notes and rubric

Design a Real-Time Payment Fraud Detection Pipeline

Problem

A payment processor handles card transactions. Before authorising, it must decide approve / review / decline within the authorisation latency budget. Design the data system that provides real-time features, runs rules and an ML model, alerts analysts, and builds training data.


Clarifying questions

QuestionAssumed answer
Throughput?10k TPS average, 50k TPS peak
Latency budget for the fraud decision?p99 < 100 ms (synchronous, in the auth path)
Fail open or closed if fraud service is down?Fail open for low amounts, rules-only fallback for high amounts
Labels?Chargebacks arrive 30–90 days later; analyst decisions within hours
Retention?7 years (compliance)
Explainability?Required: reason codes for declines

1. Requirements

2. Estimates

Peak 50k TPS × 1 KB = 50 MB/s into Kafka → 64–128 partitions keyed by card_id
Feature lookups: ~30 features per txn → 1.5M KV reads/s at peak → Redis cluster / DynamoDB (sharded)
Feature state: 500M active cards × ~200 B of counters ≈ 100 GB → fits in a Redis cluster
History: 10k × 86,400 ≈ 0.9B txns/day × 2 KB (txn + features + decision) ≈ 1.8 TB/day raw → ~300 GB/day Parquet → 7 yrs ≈ 750 TB (tiered storage)

3. Architecture

The key insight: separate the synchronous scoring path from asynchronous feature computation.

flowchart LR
    subgraph Sync["Synchronous path (p99 < 100 ms)"]
        AUTH[Auth service] -->|txn| FS[Fraud scoring service]
        FS -->|get features| ON[(Online feature store<br/>Redis / DynamoDB)]
        FS --> RULES[Rules engine<br/>hot-reloaded]
        FS --> MODEL[Model server<br/>GBDT, < 10 ms]
        FS -->|decision + reasons| AUTH
    end
    FS -->|txn + features + decision| K[[Kafka: scored_txns]]
    AUTH -->|txn events| K2[[Kafka: txns]]
    subgraph Async["Asynchronous path"]
        K2 --> FL[Flink: streaming features<br/>keyed by card_id, windows]
        FL --> ON
        K --> LAKE[(Delta: decisions + features<br/>audit log)]
        K --> CASE[Case management<br/>analyst queue]
        LBL[Chargebacks / analyst labels] --> LAKE
        LAKE --> BATCH[Batch features<br/>30/90-day aggregates]
        BATCH --> ON
        LAKE --> TRAIN[Training: point-in-time sets]
        TRAIN --> REG[Model registry] --> MODEL
    end

4. Data model

Table / storeGrainNotes
txns topic1 per transactionkey = card_id
Online storekey = card:{id}, merchant:{id}, device:{fp}latest feature values, TTL
silver.scored_txns1 per txntxn fields, feature vector used, model version, rule hits, score, decision, latency
silver.labels1 per txn label eventsource (chargeback/analyst), label_ts
gold.training_settxn × labelpoint-in-time features + label, snapshot per model version

Logging the exact feature values used at decision time is critical: it makes training data point-in-time correct by construction, and it’s your audit trail (“why was this declined?“).

5. Deep dives

txns.keyBy(t -> t.cardId)
    .process(new KeyedProcessFunction<String, Txn, FeatureUpdate>() {
        // state: ring buffer of last N txns (ts, amount, merchant, geo) with 24h TTL
        // on each txn: update counts for 1m/10m/1h/24h windows, distinct merchants (HLL),
        // last location; emit updated feature row → sink to Redis
    });

5.2 Scoring service latency budget

network in/out          10 ms
feature fetch (batched MGET, parallel stores)  10 ms p99
rules evaluation         2 ms
model inference (GBDT, ~200 trees, local in-process)  5 ms
decision + logging (async to Kafka)  2 ms
headroom                ~70 ms

5.3 Rules engine + ML combination

flowchart LR
    T[Txn + features] --> HR{Hard rules<br/>blocklist, impossible travel,<br/>velocity > X}
    HR -->|hit| DEC[Decline + reason]
    HR -->|no hit| M[ML score 0–1]
    M --> TH{Thresholds by<br/>amount / merchant risk}
    TH -->|"> 0.95"| DEC
    TH -->|"0.7–0.95"| REV[Approve + review queue<br/>or step-up auth 3DS]
    TH -->|"< 0.7"| APP[Approve]

Rules stored as versioned config (DSL/JSON), hot-reloaded, every rule hit logged with rule version, so analysts can ship a rule for a new fraud pattern in minutes without retraining.

5.4 Training data and labels

6. Trade-offs

DecisionChoiceWhy / cost
Stream engine for featuresFlinkms latency, rich keyed state + timers; Spark SS micro-batch adds seconds of feature lag
Online storeRedis cluster (or DynamoDB)sub-ms reads; cost of memory; DynamoDB if you want zero ops and multi-region
Model typeGradient boosted treesFast CPU inference, strong on tabular, explainable (SHAP reason codes)
Sync vs async scoringSync for decision, async for enrichmentLatency budget; async deep-model scoring can trigger post-auth review
Fail-open vs closedTiered by amountBusiness trade-off: lost revenue vs fraud loss. Make it explicit and configurable

7. Failure modes

FailureImpactMitigation
Feature store slow/downLatency breachTimeouts → defaults; rules-only mode; multi-AZ replicas
Flink lagStale velocity features → fraud slipsLag alerting; rules on server-side counts in scoring service as backup
Model server bad deployMass declinesShadow + canary, automatic rollback on decline-rate anomaly
Kafka unavailable for loggingMissing audit trailLocal buffer/outbox in scoring service; never block auth on logging
Fraud pattern shiftModel driftMonitor score distribution, decline rate, chargeback rate by segment; retrain cadence

8. Scaling 10×

Shard feature store by card_id; regional deployments (data residency + latency) with per-region Kafka/Flink; precompute merchant/device features in batch; GPU or optimized model runtimes only if model complexity grows.

9. What separates a senior answer

10. Follow-up questions

How do you detect "impossible travel"?

Online store keeps last transaction location + timestamp per card. On a new txn: distance (haversine) / time delta > ~900 km/h → flag. Edge cases: card-not-present transactions (merchant location ≠ cardholder), VPNs, airports; so feed it as a feature/rule with moderate weight, not an automatic decline.

How do you backfill a new feature for training?

Compute it historically in batch from the transaction history with AS-OF semantics (only data before each txn timestamp, e.g. with window functions RANGE BETWEEN INTERVAL 24 HOURS PRECEDING AND 1 MICROSECOND PRECEDING). Validate that the batch definition matches the streaming one on a recent overlap period (parity check) before training on it.

Multi-region active-active: what's hard?

Card velocity features need a global view, but a card is usually used in one region at a time. Route by card home region or replicate feature updates cross-region asynchronously and accept slight staleness; decide consistency per feature. Kafka cluster linking for logs; lake per region with residency rules.


Self-assessment rubric