Skip to main content

Building a Real-Time Fraud Detection Pipeline with Gnok

· 10 min read
Gnok Team

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:

ComponentGnok FeaturePurpose
Data ingestionCOPY INTOBulk load historical transactions
User profilingCTAS + aggregatesCompute behavioral baselines
Amount anomalyZ-score calculationFlag unusual transaction amounts
Velocity checkWindowed countsDetect burst transaction patterns
Geographic checkLEFT JOIN anti-patternFlag new-country transactions
Merchant riskCOSINE_DISTANCE + VECTOR_AVGProfile merchant similarity
Combined scoringWeighted CASE expressionUnified 0-100 risk score
DashboardPostgreSQL-compatible endpointBI visualization
InvestigationASK (NLI)Natural language forensics
ForensicsTime travelHistorical state reconstruction
MaintenanceOPTIMIZE + VACUUMKeep 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​