374 lines
No EOL
9.4 KiB
Python
374 lines
No EOL
9.4 KiB
Python
import os
|
|
import requests
|
|
import pandas as pd
|
|
import mlflow
|
|
import mlflow.data
|
|
import mlflow.pyfunc
|
|
from mlflow.models.signature import infer_signature
|
|
import pickle
|
|
import tempfile
|
|
import numpy as np
|
|
import torch
|
|
import torch.nn as nn
|
|
from torch.utils.data import DataLoader, TensorDataset
|
|
import joblib# Configuration
|
|
PERSONAL_API_KEY = os.getenv("PERSONAL_API_KEY")
|
|
POSTHOG_PROJECT_ID = os.getenv("POSTHOG_PROJECT_ID")
|
|
POSTHOG_HOST = os.getenv("POSTHOG_HOST")
|
|
|
|
def load_query_template(filepath: str, start_date: str, end_date: str) -> str:
|
|
"""Reads SQL file and replaces {start_date}/{end_date} with formatted strings."""
|
|
with open(filepath, 'r') as f:
|
|
template = f.read()
|
|
return template.format(start_date=start_date, end_date=end_date)
|
|
|
|
def fetch_posthog_data(query_string: str) -> pd.DataFrame:
|
|
"""Executes HogQL/ClickHouse query on the PostHog API."""
|
|
response = requests.post(
|
|
f"{POSTHOG_HOST}/api/projects/{POSTHOG_PROJECT_ID}/query/",
|
|
headers={
|
|
"Authorization": f"Bearer {PERSONAL_API_KEY}",
|
|
"Content-Type": "application/json",
|
|
},
|
|
json={
|
|
"query": {"kind": "HogQLQuery", "query": query_string},
|
|
},
|
|
)
|
|
response.raise_for_status()
|
|
data = response.json()
|
|
return pd.DataFrame(data=data['results'], columns=data['columns'])
|
|
|
|
class Autoencoder(nn.Module):
|
|
def __init__(self, input_dim):
|
|
super().__init__()
|
|
|
|
self.encoder = nn.Sequential(
|
|
nn.Linear(input_dim, 64),
|
|
nn.ReLU(),
|
|
nn.Linear(64, 32),
|
|
nn.ReLU(),
|
|
nn.Linear(32, 16),
|
|
nn.ReLU(),
|
|
)
|
|
|
|
self.decoder = nn.Sequential(
|
|
nn.Linear(16, 32),
|
|
nn.ReLU(),
|
|
nn.Linear(32, 64),
|
|
nn.ReLU(),
|
|
nn.Linear(64, input_dim),
|
|
)
|
|
|
|
def forward(self, x):
|
|
z = self.encoder(x)
|
|
return self.decoder(z)
|
|
|
|
from sklearn.preprocessing import StandardScaler
|
|
from sklearn.model_selection import train_test_split
|
|
|
|
|
|
class AutoencoderWithScores(mlflow.pyfunc.PythonModel):
|
|
def load_context(self, context):
|
|
self.scaler = joblib.load(context.artifacts["scaler"])
|
|
|
|
with open(context.artifacts["input_dim"], "r") as f:
|
|
input_dim = int(f.read())
|
|
|
|
self.model = Autoencoder(input_dim)
|
|
self.model.load_state_dict(
|
|
torch.load(context.artifacts["weights"], map_location="cpu")
|
|
)
|
|
self.model.eval()
|
|
|
|
with open(context.artifacts["threshold"], "r") as f:
|
|
self.threshold = float(f.read())
|
|
|
|
def predict(self, context, model_input):
|
|
X = self.scaler.transform(model_input)
|
|
X = X.astype(np.float32)
|
|
|
|
with torch.inference_mode():
|
|
reconstructed = self.model(torch.from_numpy(X)).numpy()
|
|
|
|
reconstruction_error = np.mean(
|
|
(X - reconstructed) ** 2,
|
|
axis=1,
|
|
)
|
|
|
|
predictions = np.where(
|
|
reconstruction_error > self.threshold,
|
|
-1,
|
|
1,
|
|
)
|
|
|
|
return pd.DataFrame(
|
|
{
|
|
"prediction": predictions,
|
|
"anomaly_score": reconstruction_error,
|
|
"threshold": self.threshold,
|
|
}
|
|
)
|
|
|
|
|
|
def run_training_pipeline():
|
|
mlflow.set_tracking_uri("http://127.0.0.1:8000")
|
|
mlflow.set_experiment("Session_Fraud_Detection")
|
|
|
|
sql_filepath = "./fetch_clean_sessions.sql"
|
|
|
|
start_date = "2026-03-21 00:00:00"
|
|
end_date = "2026-04-21 00:00:00"
|
|
|
|
populated_query = load_query_template(
|
|
sql_filepath,
|
|
start_date=start_date,
|
|
end_date=end_date,
|
|
)
|
|
|
|
print("Fetching training baseline from PostHog...")
|
|
|
|
df_raw = fetch_posthog_data(populated_query)
|
|
|
|
print(len(df_raw), "rows fetched from PostHog.")
|
|
|
|
start_date = "2026-04-21 00:00:01"
|
|
end_date = "2026-05-21 00:00:00"
|
|
|
|
populated_query = load_query_template(
|
|
sql_filepath,
|
|
start_date=start_date,
|
|
end_date=end_date,
|
|
)
|
|
|
|
df_raw_2 = fetch_posthog_data(populated_query)
|
|
|
|
print(len(df_raw_2), "rows fetched from PostHog.")
|
|
|
|
df_raw = pd.concat(
|
|
[df_raw, df_raw_2],
|
|
ignore_index=True,
|
|
)
|
|
|
|
feature_cols = [
|
|
c
|
|
for c in df_raw.columns
|
|
if c not in [
|
|
"person_id",
|
|
"session_id",
|
|
"session_start",
|
|
"session_end",
|
|
]
|
|
]
|
|
|
|
df_features = df_raw[feature_cols].astype(float)
|
|
|
|
with mlflow.start_run(run_name="PostHog_Autoencoder_Training"):
|
|
|
|
scaler = StandardScaler()
|
|
|
|
X = scaler.fit_transform(df_features).astype(np.float32)
|
|
|
|
X_train, X_val = train_test_split(
|
|
X,
|
|
test_size=0.2,
|
|
random_state=42,
|
|
)
|
|
|
|
train_loader = DataLoader(
|
|
TensorDataset(torch.from_numpy(X_train)),
|
|
batch_size=256,
|
|
shuffle=True,
|
|
)
|
|
|
|
model = Autoencoder(X.shape[1])
|
|
|
|
optimizer = torch.optim.Adam(
|
|
model.parameters(),
|
|
lr=1e-3,
|
|
)
|
|
|
|
criterion = nn.MSELoss()
|
|
|
|
best_loss = float("inf")
|
|
patience = 10
|
|
patience_counter = 0
|
|
|
|
for epoch in range(100):
|
|
|
|
model.train()
|
|
|
|
train_loss = 0.0
|
|
|
|
for (batch,) in train_loader:
|
|
|
|
optimizer.zero_grad()
|
|
|
|
reconstructed = model(batch)
|
|
|
|
loss = criterion(reconstructed, batch)
|
|
|
|
loss.backward()
|
|
|
|
optimizer.step()
|
|
|
|
train_loss += loss.item()
|
|
|
|
train_loss /= len(train_loader)
|
|
|
|
model.eval()
|
|
|
|
with torch.inference_mode():
|
|
|
|
val_tensor = torch.from_numpy(X_val)
|
|
|
|
val_reconstruction = model(val_tensor)
|
|
|
|
val_loss = criterion(
|
|
val_reconstruction,
|
|
val_tensor,
|
|
).item()
|
|
|
|
print(
|
|
f"Epoch {epoch+1} | "
|
|
f"Train: {train_loss:.6f} | "
|
|
f"Val: {val_loss:.6f}"
|
|
)
|
|
|
|
if val_loss < best_loss:
|
|
|
|
best_loss = val_loss
|
|
|
|
best_state = {
|
|
k: v.cpu().clone()
|
|
for k, v in model.state_dict().items()
|
|
}
|
|
|
|
patience_counter = 0
|
|
|
|
else:
|
|
|
|
patience_counter += 1
|
|
|
|
if patience_counter >= patience:
|
|
print("Early stopping.")
|
|
break
|
|
|
|
model.load_state_dict(best_state)
|
|
|
|
model.eval()
|
|
|
|
with torch.inference_mode():
|
|
|
|
val_reconstructed = model(
|
|
torch.from_numpy(X_val)
|
|
).numpy()
|
|
|
|
reconstruction_error = np.mean(
|
|
(X_val - val_reconstructed) ** 2,
|
|
axis=1,
|
|
)
|
|
|
|
threshold = np.percentile(
|
|
reconstruction_error,
|
|
99,
|
|
)
|
|
|
|
example_input = df_features.head(1)
|
|
|
|
example_scaled = scaler.transform(example_input).astype(np.float32)
|
|
|
|
with torch.inference_mode():
|
|
|
|
reconstructed = model(
|
|
torch.from_numpy(example_scaled)
|
|
).numpy()
|
|
|
|
example_error = np.mean(
|
|
(example_scaled - reconstructed) ** 2,
|
|
axis=1,
|
|
)
|
|
|
|
example_output = pd.DataFrame(
|
|
{
|
|
"prediction": np.where(
|
|
example_error > threshold,
|
|
-1,
|
|
1,
|
|
),
|
|
"anomaly_score": example_error,
|
|
"threshold": threshold,
|
|
}
|
|
)
|
|
|
|
signature = infer_signature(
|
|
example_input,
|
|
example_output,
|
|
)
|
|
|
|
mlflow.log_params(
|
|
{
|
|
"architecture": "64-32-16-32-64",
|
|
"optimizer": "Adam",
|
|
"learning_rate": 1e-3,
|
|
"epochs": epoch + 1,
|
|
"batch_size": 256,
|
|
"threshold_percentile": 99,
|
|
}
|
|
)
|
|
|
|
with tempfile.TemporaryDirectory() as tmpdir:
|
|
|
|
scaler_path = os.path.join(
|
|
tmpdir,
|
|
"scaler.pkl",
|
|
)
|
|
|
|
weights_path = os.path.join(
|
|
tmpdir,
|
|
"model.pt",
|
|
)
|
|
|
|
threshold_path = os.path.join(
|
|
tmpdir,
|
|
"threshold.txt",
|
|
)
|
|
|
|
input_dim_path = os.path.join(
|
|
tmpdir,
|
|
"input_dim.txt",
|
|
)
|
|
|
|
joblib.dump(
|
|
scaler,
|
|
scaler_path,
|
|
)
|
|
|
|
torch.save(
|
|
model.state_dict(),
|
|
weights_path,
|
|
)
|
|
|
|
with open(threshold_path, "w") as f:
|
|
f.write(str(threshold))
|
|
|
|
with open(input_dim_path, "w") as f:
|
|
f.write(str(X.shape[1]))
|
|
|
|
mlflow.pyfunc.log_model(
|
|
artifact_path="autoencoder_model",
|
|
python_model=AutoencoderWithScores(),
|
|
artifacts={
|
|
"weights": weights_path,
|
|
"scaler": scaler_path,
|
|
"threshold": threshold_path,
|
|
"input_dim": input_dim_path,
|
|
},
|
|
signature=signature,
|
|
registered_model_name="Autoencoder_Anomaly_Detector",
|
|
)
|
|
|
|
print("Model tracking and registry entry complete.")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
run_training_pipeline() |