-- Ensure clean setup if restarting from scratch DROP TABLE IF EXISTS session_predictions CASCADE; DROP TABLE IF EXISTS batch_runs CASCADE; -- ============================================================ -- BATCH RUN METADATA + AGGREGATED STATISTICS -- ============================================================ CREATE TABLE batch_runs ( batch_id UUID PRIMARY KEY DEFAULT gen_random_uuid(), -- Airflow metadata dag_run_id VARCHAR(255) NOT NULL UNIQUE, -- Time window processed by this batch window_start TIMESTAMP WITH TIME ZONE NOT NULL, window_end TIMESTAMP WITH TIME ZONE NOT NULL, -- Model lineage mlflow_model_version VARCHAR(50) NOT NULL, -- Execution status status VARCHAR(20) NOT NULL CHECK (status IN ('RUNNING', 'SUCCESS', 'FAILED')), -- Batch statistics total_sessions_processed INTEGER DEFAULT 0, total_anomalies_detected INTEGER DEFAULT 0, contamination_rate DOUBLE PRECISION, mean_anomaly_score DOUBLE PRECISION, max_anomaly_score DOUBLE PRECISION, -- Dataset drift monitoring feature_means_summary JSONB, created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP, CONSTRAINT unique_window UNIQUE (window_start, window_end) ); CREATE INDEX idx_batch_window ON batch_runs(window_start, window_end); CREATE INDEX idx_batch_created_at ON batch_runs(created_at); -- ============================================================ -- SESSION-LEVEL PREDICTIONS -- ============================================================ CREATE TABLE session_predictions ( prediction_id UUID PRIMARY KEY DEFAULT gen_random_uuid(), batch_id UUID NOT NULL REFERENCES batch_runs(batch_id) ON DELETE CASCADE, session_id VARCHAR(255) NOT NULL, user_id VARCHAR(255), session_start_time TIMESTAMP WITH TIME ZONE NOT NULL, -- Model outputs anomaly_score DOUBLE PRECISION NOT NULL, is_anomaly BOOLEAN NOT NULL, -- Dynamic payloads session_features JSONB NOT NULL, shap_values JSONB NOT NULL, created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP, CONSTRAINT unique_session_per_batch UNIQUE(batch_id, session_id) ); -- Frequently queried anomaly rows CREATE INDEX idx_session_predictions_anomalies ON session_predictions(batch_id) WHERE is_anomaly = TRUE; -- Optional searches inside feature JSON CREATE INDEX idx_session_features_gin ON session_predictions USING GIN(session_features);