Skip to main content

Streaming ML

Streaming inference applies an existing model to new events. Online learning updates a model from labelled outcomes. Anomaly detection compares incoming measurements with a baseline. These solve different problems: scoring transactions does not train a fraud model, and an anomaly score is not a fraud probability.

Prerequisites​

Start with the streaming tutorial for source tables, model training, and a controlled sequence of input rows. Derived streams require an authenticated tenant and a source in the same catalog/schema. Keep the catalog and background services available and grant the managed tenant service identity the access its workload needs. See setup.

Streaming inference​

First register a source stream over a table, then define the derived output. These fragments use the tutorial's existing tables and trained model:

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;

The scaling matches model training. A SELECT-based definition can also filter and project events:

CREATE OR REPLACE STREAM tutorial.streaming.rt_scored AS
SELECT txn_id, 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 the source before inserting the events you want it to consume. APPEND_ONLY does not represent update/delete change processing.

Processing and inspection​

REFRESH STREAM tutorial.streaming.fraud_predictions;
REFRESH STREAM tutorial.streaming.rt_scored;
SHOW STREAMS ON TABLE tutorial.streaming.transactions;

SELECT txn_id, amount, p_fraud
FROM tutorial.streaming.fraud_predictions ORDER BY txn_id;

A successful definition is not evidence that events were processed. Confirm expected row IDs and compare stream predictions with direct calls to the same active model using identical features. The tutorial includes this comparison.

Anomaly detection​

The durable anomaly SQL path currently supports METHOD 'ema':

CREATE OR REPLACE STREAM tutorial.streaming.anomalies AS
DETECT ANOMALIES IN tutorial.streaming.txn_stream
ON amount USING METHOD 'ema' WINDOW '5m';

REFRESH STREAM tutorial.streaming.anomalies;
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;

The prior-baseline columns help explain the result. Feed representative normal events before evaluating anomalous probes; warm-up rows are not evidence of a trained detector. See anomaly detection for interpretation and library-only methods.

WINDOW accepts a positive integer with an ms, s, m, or h suffix, up to 24 hours. Do not infer event-time ordering or late-event behavior merely from the presence of a window clause.

Online learning​

Native linear and binary logistic models support incremental updates from labelled feedback. After preparing the online-learning tutorial:

ALTER MODEL tutorial.online.scorer
ENABLE ONLINE LEARNING
WITH FEEDBACK FROM tutorial.online.labeled_outcomes
LEARNING RATE 0.01;

SHOW ONLINE LEARNING;

An enabled learner with an empty feedback table has no new observations to learn from. Confirm consumed feedback, run outcomes, and model changes after inserting labelled events. Keep the feedback feature order, types, and preprocessing consistent with the original training query.

Slices let you observe configured subpopulations:

ALTER MODEL tutorial.online.scorer
ENABLE ONLINE LEARNING SLICE AS 'eu_customers'
WHERE region = 'eu';

SHOW MODEL SLICES tutorial.online.scorer;
SHOW AUTOML RUNS LIMIT 20;

Use separate held-out data to check whether updates improve quality. Automatic safeguards do not replace domain-specific acceptance criteria.

ALTER MODEL tutorial.online.scorer DISABLE ONLINE LEARNING;

Operational checks​

Inspect progress and failures after a restart or dependency outage. Check source access, active model version, pending data, and output rows before retrying a workload. Preserve stable event identifiers so missing or repeated results can be detected by the application.

Streaming tutorial · Online-learning tutorial · Feature store