Building a Real-Time Fraud Detection Pipeline with Gnok
This tutorial builds a complete fraud detection system inside Gnok -- from data ingestion to ML scoring to alerting -- without any external services. We'll use COPY INTO for bulk loading, vector similarity for merchant profiling, statistical anomaly detection for flagging outliers, and natural language queries for ad-hoc investigation.
Architecture
Everything runs inside Gnok. No Kafka, no external model servers, no ETL pipelines.
Step 1: Set Up the Schema
CREATE SCHEMA fraud_demo;
USE SCHEMA fraud_demo;
-- Transaction table
CREATE TABLE transactions (
tx_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
merchant_id BIGINT,
amount DECIMAL(10, 2),
currency VARCHAR,
tx_type VARCHAR,
channel VARCHAR,
country VARCHAR,
is_international BOOLEAN,
hour_of_day INT,
day_of_week INT,
tx_timestamp TIMESTAMP
)
PARTITIONED BY (day(tx_timestamp));
-- Merchant table with embeddings
CREATE TABLE merchants (
merchant_id BIGINT NOT NULL,
name VARCHAR,
category VARCHAR,
risk_tier VARCHAR,
avg_ticket DECIMAL(10, 2),
embedding VECTOR(128)
);
Step 2: Bulk Load Historical Data
Upload transaction data to a stage, then load it:
-- Load historical transactions from Parquet
COPY INTO transactions
FROM '@stage/fraud_demo/transactions/'
FILE_FORMAT = (TYPE = PARQUET);
-- Load merchant data from CSV
COPY INTO merchants
FROM '@stage/fraud_demo/merchants/'
FILE_FORMAT = (
TYPE = CSV,
DELIMITER = ',',
SKIP_ROWS = 1
);
Gnok's COPY INTO automatically parallelizes across workers and routes rows to the correct partitions. For our transaction table partitioned by day(tx_timestamp), each day's data ends up in optimally-sized files.
Verify the Load
SELECT COUNT(*) AS total_txns,
MIN(tx_timestamp) AS earliest,
MAX(tx_timestamp) AS latest
FROM transactions;
SELECT COUNT(*) AS merchant_count,
COUNT(DISTINCT category) AS categories
FROM merchants;
Step 3: Build User Behavior Profiles
Fraud detection works best when you can compare a transaction against the user's normal behavior. Let's compute per-user statistics:
CREATE TABLE user_profiles AS
SELECT
user_id,
COUNT(*) AS tx_count_total,
AVG(amount) AS avg_amount,
STDDEV(amount) AS stddev_amount,
PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY amount) AS p95_amount,
COUNT(DISTINCT merchant_id) AS unique_merchants,
COUNT(DISTINCT country) AS unique_countries,
SUM(CASE WHEN is_international THEN 1 ELSE 0 END) AS intl_tx_count,
MODE(channel) AS preferred_channel,
MAX(tx_timestamp) AS last_tx_time
FROM transactions
GROUP BY user_id;
Check the profile for a specific user:
SELECT * FROM user_profiles WHERE user_id = 12345;
Step 4: Statistical Anomaly Scoring
Now let's build a scoring query that flags anomalies using multiple signals.
Amount Anomaly (Z-Score)
Compare each transaction's amount against the user's historical distribution:
SELECT
t.tx_id,
t.user_id,
t.amount,
u.avg_amount,
u.stddev_amount,
-- Z-score: how many standard deviations from the user's mean
CASE
WHEN u.stddev_amount > 0
THEN (t.amount - u.avg_amount) / u.stddev_amount
ELSE 0
END AS amount_zscore,
-- Flag if more than 2.5 standard deviations
CASE
WHEN u.stddev_amount > 0
AND ABS((t.amount - u.avg_amount) / u.stddev_amount) > 2.5
THEN true
ELSE false
END AS amount_anomaly
FROM transactions t
JOIN user_profiles u ON t.user_id = u.user_id
WHERE t.tx_timestamp > CURRENT_TIMESTAMP - INTERVAL '1' HOUR;
Velocity Anomaly
Flag users with unusual transaction frequency:
WITH recent_velocity AS (
SELECT
user_id,
COUNT(*) AS tx_last_hour,
SUM(amount) AS total_last_hour
FROM transactions
WHERE tx_timestamp > CURRENT_TIMESTAMP - INTERVAL '1' HOUR
GROUP BY user_id
)
SELECT
v.user_id,
v.tx_last_hour,
v.total_last_hour,
-- Compare against historical daily average
u.tx_count_total / 30.0 AS avg_daily_rate,
-- Velocity ratio: how many times above normal
v.tx_last_hour / GREATEST(u.tx_count_total / 720.0, 0.1) AS velocity_ratio
FROM recent_velocity v
JOIN user_profiles u ON v.user_id = u.user_id
WHERE v.tx_last_hour / GREATEST(u.tx_count_total / 720.0, 0.1) > 5.0;
Geographic Anomaly
Flag transactions from countries the user has never used before:
WITH user_countries AS (
SELECT user_id, country
FROM transactions
GROUP BY user_id, country
)
SELECT
t.tx_id,
t.user_id,
t.country AS tx_country,
true AS new_country_flag
FROM transactions t
LEFT JOIN user_countries uc
ON t.user_id = uc.user_id AND t.country = uc.country
WHERE t.tx_timestamp > CURRENT_TIMESTAMP - INTERVAL '1' HOUR
AND uc.country IS NULL;
Step 5: Merchant Similarity Analysis
Use vector distance to find merchants similar to known high-risk merchants:
-- Find merchants whose embeddings are closest to a known fraudulent merchant
WITH fraud_merchant AS (
SELECT embedding
FROM merchants
WHERE merchant_id = 9999 -- known fraud merchant
)
SELECT
m.merchant_id,
m.name,
m.category,
COSINE_DISTANCE(m.embedding, fm.embedding) AS similarity_distance
FROM merchants m
CROSS JOIN fraud_merchant fm
WHERE m.merchant_id <> 9999
ORDER BY similarity_distance ASC
LIMIT 20;
Compute average embedding per merchant category for cluster analysis:
SELECT
category,
COUNT(*) AS merchant_count,
VECTOR_AVG(embedding) AS category_centroid
FROM merchants
GROUP BY category;
Step 6: Combined Risk Score
Pull all the signals together into a single scoring query:
CREATE VIEW fraud_scores AS
WITH amount_scores AS (
SELECT
t.tx_id,
t.user_id,
t.merchant_id,
t.amount,
t.country,
t.tx_timestamp,
t.is_international,
CASE
WHEN u.stddev_amount > 0
THEN ABS((t.amount - u.avg_amount) / u.stddev_amount)
ELSE 0
END AS amount_zscore,
t.amount > u.p95_amount AS exceeds_p95
FROM transactions t
JOIN user_profiles u ON t.user_id = u.user_id
WHERE t.tx_timestamp > CURRENT_TIMESTAMP - INTERVAL '1' HOUR
),
velocity_scores AS (
SELECT
user_id,
COUNT(*) AS tx_last_hour
FROM transactions
WHERE tx_timestamp > CURRENT_TIMESTAMP - INTERVAL '1' HOUR
GROUP BY user_id
),
geo_scores AS (
SELECT DISTINCT t.tx_id, true AS new_country
FROM transactions t
LEFT JOIN (
SELECT user_id, country FROM transactions GROUP BY user_id, country
) hist ON t.user_id = hist.user_id AND t.country = hist.country
WHERE t.tx_timestamp > CURRENT_TIMESTAMP - INTERVAL '1' HOUR
AND hist.country IS NULL
)
SELECT
a.tx_id,
a.user_id,
a.amount,
a.country,
a.tx_timestamp,
m.category AS merchant_category,
m.risk_tier AS merchant_risk_tier,
-- Individual signals
a.amount_zscore,
a.exceeds_p95,
COALESCE(v.tx_last_hour, 0) AS tx_last_hour,
COALESCE(g.new_country, false) AS new_country,
-- Combined risk score (weighted sum, 0-100)
LEAST(100, ROUND(
CASE WHEN a.amount_zscore > 3.0 THEN 30 WHEN a.amount_zscore > 2.0 THEN 15 ELSE 0 END
+ CASE WHEN a.exceeds_p95 THEN 15 ELSE 0 END
+ CASE WHEN COALESCE(v.tx_last_hour, 0) > 10 THEN 25 WHEN COALESCE(v.tx_last_hour, 0) > 5 THEN 10 ELSE 0 END
+ CASE WHEN COALESCE(g.new_country, false) THEN 20 ELSE 0 END
+ CASE WHEN m.risk_tier = 'high' THEN 15 WHEN m.risk_tier = 'medium' THEN 5 ELSE 0 END
+ CASE WHEN a.is_international THEN 5 ELSE 0 END
)) AS risk_score,
-- Risk level classification
CASE
WHEN LEAST(100, ROUND(
CASE WHEN a.amount_zscore > 3.0 THEN 30 WHEN a.amount_zscore > 2.0 THEN 15 ELSE 0 END
+ CASE WHEN a.exceeds_p95 THEN 15 ELSE 0 END
+ CASE WHEN COALESCE(v.tx_last_hour, 0) > 10 THEN 25 WHEN COALESCE(v.tx_last_hour, 0) > 5 THEN 10 ELSE 0 END
+ CASE WHEN COALESCE(g.new_country, false) THEN 20 ELSE 0 END
+ CASE WHEN m.risk_tier = 'high' THEN 15 WHEN m.risk_tier = 'medium' THEN 5 ELSE 0 END
+ CASE WHEN a.is_international THEN 5 ELSE 0 END
)) >= 60 THEN 'CRITICAL'
WHEN LEAST(100, ROUND(
CASE WHEN a.amount_zscore > 3.0 THEN 30 WHEN a.amount_zscore > 2.0 THEN 15 ELSE 0 END
+ CASE WHEN a.exceeds_p95 THEN 15 ELSE 0 END
+ CASE WHEN COALESCE(v.tx_last_hour, 0) > 10 THEN 25 WHEN COALESCE(v.tx_last_hour, 0) > 5 THEN 10 ELSE 0 END
+ CASE WHEN COALESCE(g.new_country, false) THEN 20 ELSE 0 END
+ CASE WHEN m.risk_tier = 'high' THEN 15 WHEN m.risk_tier = 'medium' THEN 5 ELSE 0 END
+ CASE WHEN a.is_international THEN 5 ELSE 0 END
)) >= 40 THEN 'HIGH'
WHEN LEAST(100, ROUND(
CASE WHEN a.amount_zscore > 3.0 THEN 30 WHEN a.amount_zscore > 2.0 THEN 15 ELSE 0 END
+ CASE WHEN a.exceeds_p95 THEN 15 ELSE 0 END
+ CASE WHEN COALESCE(v.tx_last_hour, 0) > 10 THEN 25 WHEN COALESCE(v.tx_last_hour, 0) > 5 THEN 10 ELSE 0 END
+ CASE WHEN COALESCE(g.new_country, false) THEN 20 ELSE 0 END
+ CASE WHEN m.risk_tier = 'high' THEN 15 WHEN m.risk_tier = 'medium' THEN 5 ELSE 0 END
+ CASE WHEN a.is_international THEN 5 ELSE 0 END
)) >= 20 THEN 'MEDIUM'
ELSE 'LOW'
END AS risk_level
FROM amount_scores a
LEFT JOIN velocity_scores v ON a.user_id = v.user_id
LEFT JOIN geo_scores g ON a.tx_id = g.tx_id
LEFT JOIN merchants m ON a.merchant_id = m.merchant_id;
Query the Scores
-- Top 20 highest-risk transactions in the last hour
SELECT * FROM fraud_scores
ORDER BY risk_score DESC
LIMIT 20;
-- Summary by risk level
SELECT
risk_level,
COUNT(*) AS tx_count,
SUM(amount) AS total_amount,
AVG(risk_score) AS avg_score
FROM fraud_scores
GROUP BY risk_level
ORDER BY avg_score DESC;
-- Critical alerts with full context
SELECT
tx_id, user_id, amount, country,
merchant_category, merchant_risk_tier,
risk_score,
amount_zscore,
tx_last_hour,
new_country
FROM fraud_scores
WHERE risk_level = 'CRITICAL'
ORDER BY risk_score DESC;
Step 7: Connect a Dashboard
Connect your BI tool to your account's assigned PostgreSQL-compatible endpoint (see Connect an Application and PostgreSQL wire protocol), or build a Studio dashboard, on the fraud_scores view:
Dashboard panels:
- Risk distribution — pie chart of risk_level counts
- Top flagged transactions — table sorted by risk_score DESC
- Risk by merchant category — bar chart of avg risk_score per category
- Hourly trend — time series of CRITICAL + HIGH transaction counts
- Geographic heat map — risk scores by country
Step 8: Investigate with Natural Language
When an alert fires, use ASK for quick investigation. ASK uses Gnok's AI provider, so an organization administrator must first enable AI Governance for your organization (see Natural Language Queries):
-- Set context
USE CATALOG my_catalog;
USE SCHEMA fraud_demo;
-- Investigate
ASK 'Show me all transactions for user 12345 in the last 24 hours, sorted by time';
ASK 'What merchant categories have the most critical-risk transactions today?';
ASK 'How does today''s high-risk transaction count compare to the last 7 days?';
ASK 'Find users with more than 10 transactions in the last hour';
Step 9: Time Travel for Forensics
When investigating a confirmed fraud case, use time travel to reconstruct what happened:
-- What did the user's profile look like before the fraud?
SELECT * FROM user_profiles
AT(TIMESTAMP => '2026-03-15 08:00:00')
WHERE user_id = 12345;
-- Compare with current profile
SELECT * FROM user_profiles
WHERE user_id = 12345;
-- Track how the user's transaction pattern changed
SELECT
DATE_TRUNC('day', tx_timestamp) AS day,
COUNT(*) AS tx_count,
SUM(amount) AS total_amount,
AVG(amount) AS avg_amount,
COUNT(DISTINCT country) AS countries
FROM transactions
WHERE user_id = 12345
AND tx_timestamp > CURRENT_TIMESTAMP - INTERVAL '30' DAY
GROUP BY DATE_TRUNC('day', tx_timestamp)
ORDER BY day;
Step 10: Ongoing Maintenance
Keep the pipeline running smoothly:
-- Refresh user profiles (run hourly or on a schedule)
CREATE TABLE user_profiles_new AS
SELECT /* same query as Step 3 */ ...;
DROP TABLE user_profiles;
ALTER TABLE user_profiles_new RENAME TO user_profiles;
-- Compact transaction data (run daily)
OPTIMIZE TABLE transactions;
-- Clean up old snapshots (run weekly)
ALTER TABLE transactions EXECUTE EXPIRE_SNAPSHOTS
SET ('older-than' = '2026-03-01T00:00:00Z');
VACUUM TABLE transactions;
Summary
What we built:
| Component | Gnok Feature | Purpose |
|---|---|---|
| Data ingestion | COPY INTO | Bulk load historical transactions |
| User profiling | CTAS + aggregates | Compute behavioral baselines |
| Amount anomaly | Z-score calculation | Flag unusual transaction amounts |
| Velocity check | Windowed counts | Detect burst transaction patterns |
| Geographic check | LEFT JOIN anti-pattern | Flag new-country transactions |
| Merchant risk | COSINE_DISTANCE + VECTOR_AVG | Profile merchant similarity |
| Combined scoring | Weighted CASE expression | Unified 0-100 risk score |
| Dashboard | PostgreSQL-compatible endpoint | BI visualization |
| Investigation | ASK (NLI) | Natural language forensics |
| Forensics | Time travel | Historical state reconstruction |
| Maintenance | OPTIMIZE + VACUUM | Keep performance optimal |
Apart from the ASK investigation step, which calls Gnok's AI provider, all of this runs inside Gnok. The scoring view executes in-engine with partition pruning, SIMD-accelerated aggregation, and distributed MPP execution.
Further Reading
- AI & ML Tutorials -- model training, ONNX inference, feature groups, streaming scoring
- Vector Operations -- distance functions reference
- Time Travel -- snapshot query syntax
- COPY INTO -- bulk loading reference
- Table Maintenance -- OPTIMIZE, VACUUM schedules
- Connect an Application -- SQL clients and BI tools