Streaming scoring and anomaly detection
Create a CDC source, materialize model predictions and EMA anomalies, and verify durable output counts.
The practical problem and model
A transaction-monitoring pipeline needs a result when a new event arrives, not just when an analyst reruns a batch query. This tutorial combines a trained classifier with a change-data-capture (CDC) source and three derived streams.
The classifier uses labelled examples to score a transaction. The exponential-moving-average (EMA) detector instead compares an amount with recent amounts, using a five-minute window. A high classifier score and an unusual amount answer different questions and should not be conflated.
The setup creates 30 baseline transactions followed by two probes, including a large amount of 9,500. Explicit refresh statements make the short walkthrough observable without waiting for a scheduler tick.
Before you run
Complete the shared setup. This walkthrough uses synthetic data and recreates its tutorial objects. Use a sandbox tenant or schema, and run the steps in order.
Studio: Tutorial: Streaming ML — Anomaly Detection + Continuous Scoring. Download the complete SQL.
The engine and catalog must stay available. The tenant background identity needs access to the source and managed output workload; see background prerequisites.
Step 1: Train the scorer with scaled inputs
Amount, velocity, and distance use different units. The model trains on scaled versions and every serving expression repeats the same transformations.
CREATE CATALOG IF NOT EXISTS tutorial;
CREATE SCHEMA IF NOT EXISTS tutorial.streaming;
CREATE OR REPLACE TABLE tutorial.streaming.transactions (
txn_id BIGINT,
amount DOUBLE,
velocity_1h DOUBLE,
distance_km DOUBLE,
arrived_at TIMESTAMP
);
CREATE OR REPLACE TABLE tutorial.streaming.fraud_history (
amount DOUBLE, velocity_1h DOUBLE, distance_km DOUBLE, is_fraud INT
);
INSERT INTO tutorial.streaming.fraud_history VALUES
(50.0, 1.0, 5.0, 0), (1200.0, 8.0, 800.0, 1),
(75.0, 2.0, 12.0, 0), (980.0, 6.0, 600.0, 1),
(30.0, 1.0, 3.0, 0), (1450.0, 9.0, 950.0, 1);
CREATE OR REPLACE MODEL tutorial.streaming.fraud_scorer
(DOUBLE, DOUBLE, DOUBLE)
RETURNS DOUBLE
TYPE 'logistic'
OPTIONS (learning_rate = 0.05, max_iters = 1000)
AS SELECT CAST(is_fraud AS DOUBLE) AS label,
amount / 1500.0 AS amount_scaled,
velocity_1h / 10.0 AS velocity_scaled,
distance_km / 1000.0 AS distance_scaled
FROM tutorial.streaming.fraud_history;
Step 2: Create the source and derived streams
One stream scores every event; one applies a SQL filter before retaining scored rows; one records EMA diagnostics. The source is append-only for this fixture.
DROP STREAM IF EXISTS tutorial.streaming.txn_stream;
CREATE STREAM tutorial.streaming.txn_stream
ON TABLE tutorial.streaming.transactions APPEND_ONLY = TRUE;
CREATE OR REPLACE STREAM tutorial.streaming.fraud_predictions
WITH MODEL tutorial.streaming.fraud_scorer
USING (amount / 1500.0, velocity_1h / 10.0, distance_km / 1000.0)
AS p_fraud
FROM tutorial.streaming.txn_stream;
CREATE OR REPLACE STREAM tutorial.streaming.rt_scored AS
SELECT
txn_id,
arrived_at,
amount,
ML_PREDICT('tutorial.streaming.fraud_scorer', amount / 1500.0, velocity_1h / 10.0, distance_km / 1000.0) AS p_fraud
FROM tutorial.streaming.txn_stream
WHERE amount >= 100.0;
CREATE OR REPLACE STREAM tutorial.streaming.anomalies AS
DETECT ANOMALIES IN tutorial.streaming.txn_stream
ON amount
USING METHOD 'ema'
WINDOW '5m';
SHOW STREAMS ON TABLE tutorial.streaming.transactions;
Step 3: Insert baseline and probe events, then refresh
Thirty baseline rows establish history. The two subsequent probes have IDs 1001 and 1002. Refresh processes available source changes into the derived outputs.
INSERT INTO tutorial.streaming.transactions VALUES
(100, 60.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(101, 80.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(102, 100.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(103, 120.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(104, 140.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(105, 60.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(106, 80.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(107, 100.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(108, 120.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(109, 140.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(110, 60.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(111, 80.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(112, 100.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(113, 120.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(114, 140.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(115, 60.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(116, 80.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(117, 100.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(118, 120.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(119, 140.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(120, 60.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(121, 80.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(122, 100.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(123, 120.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(124, 140.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(125, 60.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(126, 80.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(127, 100.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(128, 120.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(129, 140.0, 1.0, 5.0, CURRENT_TIMESTAMP);
INSERT INTO tutorial.streaming.transactions VALUES
(1001, 100.0, 1.0, 5.0, CURRENT_TIMESTAMP),
(1002, 9500.0, 9.0, 950.0, CURRENT_TIMESTAMP);
REFRESH STREAM tutorial.streaming.fraud_predictions;
REFRESH STREAM tutorial.streaming.rt_scored;
REFRESH STREAM tutorial.streaming.anomalies;
Step 4: Read every output and compare direct scoring
The final direct query provides a check against materialized predictions. Inspect anomaly status and prior counts as well as the boolean flag; early events are in warmup.
SELECT txn_id, amount, p_fraud
FROM tutorial.streaming.fraud_predictions
ORDER BY txn_id;
SELECT txn_id, arrived_at, amount, p_fraud
FROM tutorial.streaming.rt_scored
ORDER BY txn_id;
SELECT txn_id, amount, anomaly_score, is_anomaly,
ema_prior_mean, ema_prior_stddev, ema_prior_count, anomaly_status
FROM tutorial.streaming.anomalies
ORDER BY txn_id;
SELECT txn_id, amount,
tutorial.streaming.fraud_scorer(amount / 1500.0, velocity_1h / 10.0, distance_km / 1000.0) AS p_fraud
FROM tutorial.streaming.txn_stream
ORDER BY txn_id;
What to check in the results
| Output | Expected rows after this reset/run |
|---|---|
| All model predictions | 32 |
Filtered scored stream (amount >= 100) | 20 |
| Anomaly diagnostics | 32 |
| Direct source scoring | 32 |
The materialized and direct classifier scores should agree by transaction ID. The large probe is anomalous after baseline history is available; warmup rows must not be interpreted as evidence of normality. EMA state follows source processing order, which is not necessarily transaction-ID order.
Separate lifecycle checks exercised additional arrivals, concurrent processing, model-version pinning, and cold restart recovery. Those checks used more rows than this 32-row walkthrough.
Adapt it to real data
Define whether your input contains inserts only or also corrections and deletes. Choose a source-order and late-arrival policy, pin the model version used for durable predictions, and retain enough lineage to explain historical results.
Monitor processing lag, failed refreshes, output duplication, and checkpoint progress. An anomaly signal is a reason to investigate; it is not by itself proof of fraud or a safe basis for blocking a real transaction.