Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/cd-production.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ on:
paths:
- 'services/model-api/**'
- 'services/alert-bridge/**'
- 'k8s/**'

jobs:
deploy:
Expand All @@ -30,6 +31,9 @@ jobs:
- name: model-api ๋ฐฐํฌ ํ™•์ธ
run: kubectl rollout status deployment/model-api -n model-api

- name: alert-bridge kubectl apply
run: kubectl apply -f k8s/alert-bridge/deployment.yaml

- name: alert-bridge ํŒŒ๋“œ ์žฌ์‹œ์ž‘
run: kubectl rollout restart deployment/alert-bridge -n monitoring

Expand Down
22 changes: 22 additions & 0 deletions .github/workflows/deploy-churn.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
name: Deploy (churn)

on:
repository_dispatch:
types:
- deploy-churn

jobs:
deploy:
runs-on: ubuntu-latest

steps:
- name: kubectl ์„ค์ •
run: |
mkdir -p $HOME/.kube
echo "${{ secrets.KUBECONFIG_PRODUCTION }}" | base64 -d > $HOME/.kube/config

- name: model-api ํŒŒ๋“œ ์žฌ์‹œ์ž‘
run: kubectl rollout restart deployment/model-api -n model-api

- name: model-api ๋ฐฐํฌ ํ™•์ธ
run: kubectl rollout status deployment/model-api -n model-api
22 changes: 22 additions & 0 deletions .github/workflows/deploy-uplift.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
name: Deploy (uplift)

on:
repository_dispatch:
types:
- deploy-uplift

jobs:
deploy:
runs-on: ubuntu-latest

steps:
- name: kubectl ์„ค์ •
run: |
mkdir -p $HOME/.kube
echo "${{ secrets.KUBECONFIG_PRODUCTION }}" | base64 -d > $HOME/.kube/config

- name: model-api ํŒŒ๋“œ ์žฌ์‹œ์ž‘
run: kubectl rollout restart deployment/model-api -n model-api

- name: model-api ๋ฐฐํฌ ํ™•์ธ
run: kubectl rollout status deployment/model-api -n model-api
46 changes: 0 additions & 46 deletions .github/workflows/retrain-churn.yaml

This file was deleted.

35 changes: 28 additions & 7 deletions k8s/alert-bridge/deployment.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -19,12 +19,33 @@ spec:
ports:
- containerPort: 8080
env:
- name: GITHUB_TOKEN
- name: SMTP_HOST
valueFrom:
secretKeyRef:
name: github-credentials
key: token
- name: GITHUB_REPO_OWNER
value: simGPT
- name: GITHUB_REPO_NAME
value: sw-mlops
name: smtp-credentials
key: host
- name: SMTP_PORT
valueFrom:
secretKeyRef:
name: smtp-credentials
key: port
- name: SMTP_USER
valueFrom:
secretKeyRef:
name: smtp-credentials
key: user
- name: SMTP_PASSWORD
valueFrom:
secretKeyRef:
name: smtp-credentials
key: password
- name: ALERT_MAIL_FROM
valueFrom:
secretKeyRef:
name: smtp-credentials
key: mail-from
- name: ALERT_MAIL_TO
valueFrom:
secretKeyRef:
name: smtp-credentials
key: mail-to
63 changes: 42 additions & 21 deletions services/alert-bridge/main.py
Original file line number Diff line number Diff line change
@@ -1,36 +1,57 @@
import os
import httpx
import smtplib
import ssl
from datetime import datetime, timezone
from email.mime.text import MIMEText
from fastapi import FastAPI, Request

app = FastAPI(title="Alert Bridge")

GITHUB_TOKEN = os.getenv("GITHUB_TOKEN")
GITHUB_REPO_OWNER = os.getenv("GITHUB_REPO_OWNER")
GITHUB_REPO_NAME = os.getenv("GITHUB_REPO_NAME")
SMTP_HOST = os.getenv("SMTP_HOST")
SMTP_PORT = int(os.getenv("SMTP_PORT", "465"))
SMTP_USER = os.getenv("SMTP_USER")
SMTP_PASSWORD = os.getenv("SMTP_PASSWORD")
ALERT_MAIL_FROM = os.getenv("ALERT_MAIL_FROM")
ALERT_MAIL_TO = os.getenv("ALERT_MAIL_TO")

# Alertmanager์—์„œ ์•Œ๋ฆผ์„ ๋ฐ›์•„์„œ GitHub Actions๋กœ ์ „๋‹ฌํ•˜๋Š” ์—ญํ• ์„ ํ•˜๋Š” api ์—”๋“œํฌ์ธํŠธ(Alertmanager๊ฐ€ ์š”์ฒญ์„ ๋ณด๋ƒ„)
# ์•Œ๋ฆผ ๋‚ด์šฉ์„ ์ด๋ฉ”์ผ๋กœ ์ž‘์„ฑํ•˜๋Š” ํ•จ์ˆ˜
def build_email_body(alerts: list) -> str:
lines = [f"[sw-mlops ์•Œ๋ฆผ] {len(alerts)}๊ฐœ์˜ ์ด์ƒ ๊ฐ์ง€ - {datetime.now(timezone.utc).strftime('%Y-%m-%d %H:%M:%S UTC')}\n"]
for a in alerts:
labels = a.get("labels", {})
annotations = a.get("annotations", {})
lines.append(f"โ€ข Alert : {labels.get('alertname', '-')}")
lines.append(f" Severity: {labels.get('severity', '-')}")
lines.append(f" Summary : {annotations.get('summary', '-')}")
lines.append(f" Detail : {annotations.get('description', '-')}\n")
return "\n".join(lines)

# ์ด๋ฉ”์ผ ์ „์†กํ•˜๋Š” ํ•จ์ˆ˜
def send_email(subject: str, body: str):
msg = MIMEText(body, "plain", "utf-8")
msg["Subject"] = subject
msg["From"] = ALERT_MAIL_FROM
msg["To"] = ALERT_MAIL_TO

with smtplib.SMTP(SMTP_HOST, SMTP_PORT) as server:
server.starttls(context=ssl.create_default_context())
server.login(SMTP_USER, SMTP_PASSWORD)
server.sendmail(ALERT_MAIL_FROM, ALERT_MAIL_TO, msg.as_string())


# ์•Œ๋ฆผ ๋ฐ›์•„์„œ ์ด๋ฉ”์ผ ์ „์†กํ•˜๋Š” api
@app.post("/alert")
async def receive_alert(request: Request):
payload = await request.json()

# alertmanager์—์„œ ๋ฐ›์€ ์•Œ๋ฆผ ์ค‘์—์„œ firing ์ƒํƒœ์ธ ์•Œ๋ฆผ๋งŒ ํ•„ํ„ฐ๋ง
firing_alerts = [a for a in payload.get("alerts", []) if a["status"] == "firing"]
if not firing_alerts:
return {"message": "no firing alerts"}
return {"message": "๋ฐœ์ƒ์ค‘์ธ ์•Œ๋ฆผ์ด ์—†์Šต๋‹ˆ๋‹ค."}

alert_names = [a["labels"].get("alertname", "") for a in firing_alerts]
subject = f"[sw-mlops ์•Œ๋ฆผ] {', '.join(alert_names)}"
body = build_email_body(firing_alerts) # ์•Œ๋ฆผ ๋‚ด์šฉ์„ ์ด๋ฉ”์ผ body๋กœ ์ž‘์„ฑ

send_email(subject, body) # ์ด๋ฉ”์ผ ์ „์†ก

async with httpx.AsyncClient() as client:
response = await client.post(
f"https://api.github.com/repos/{GITHUB_REPO_OWNER}/{GITHUB_REPO_NAME}/dispatches",
headers={
"Authorization": f"Bearer {GITHUB_TOKEN}",
"Accept": "application/vnd.github+json",
},
json={
"event_type": "model-retrain",
"client_payload": {"alerts": alert_names},
},
)

return {"status": response.status_code, "alerts": alert_names}
return {"message": "์ด๋ฉ”์ผ์ด ์ „์†ก๋˜์—ˆ์Šต๋‹ˆ๋‹ค.", "alerts": alert_names}
1 change: 0 additions & 1 deletion services/alert-bridge/requirements.txt
Original file line number Diff line number Diff line change
@@ -1,3 +1,2 @@
fastapi
uvicorn
httpx
24 changes: 24 additions & 0 deletions services/model-api/app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,3 +61,27 @@ def churn_predict(request: RequestSchema):
metadata=None,
error=str(e),
)


# Uplift ์˜ˆ์ธก API
@app.post("/uplift_predict", response_model=ResponseSchema)
def uplift_predict(request: RequestSchema):
try:
service = get_service("uplift")
output = service.predict(request.data)

return ResponseSchema(
success=True,
model="uplift",
result=output["result"],
metadata=output["metadata"],
error=None,
)
except Exception as e:
return ResponseSchema(
success=False,
model="uplift",
result=None,
metadata=None,
error=str(e),
)
21 changes: 19 additions & 2 deletions services/model-api/app/models/loader.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
import mlflow.pytorch
import mlflow.sklearn

from app.models.uplift_model import TLearner # noqa: F401 MLflow ์—ญ์ง๋ ฌํ™” ์‹œ ํ•„์š”

_model_cache: dict = {}

Expand All @@ -11,7 +12,7 @@ def load_model(model_name: str):
if cache_key in _model_cache:
return _model_cache[cache_key]

mlflow_uri = os.getenv("MLFLOW_TRACKING_URI", "http://mlflow:5000") # ํ™˜๊ฒฝ๋ณ€์ˆ˜์—์„œ mlflow tracking uri ๊ฐ€์ ธ์˜ค๊ธฐ, ์—†์œผ๋ฉด ๊ธฐ๋ณธ๊ฐ’์œผ๋กœ http://mlflow:5000 ์‚ฌ์šฉ
mlflow_uri = os.getenv("MLFLOW_TRACKING_URI", "http://mlflow:5000") # deployment์—์„œ mlflow tracking uri ๊ฐ€์ ธ์˜ค๊ธฐ, ์—†์œผ๋ฉด ๊ธฐ๋ณธ๊ฐ’์œผ๋กœ http://mlflow:5000 ์‚ฌ์šฉ
mlflow.set_tracking_uri(mlflow_uri)

model_uri = f"models:/{model_name}/latest"
Expand All @@ -27,7 +28,23 @@ def load_churn_model(model_name: str):
if cache_key in _model_cache:
return _model_cache[cache_key]

mlflow_uri = os.getenv("MLFLOW_TRACKING_URI", "http://mlflow:5000") # ํ™˜๊ฒฝ๋ณ€์ˆ˜์—์„œ mlflow tracking uri ๊ฐ€์ ธ์˜ค๊ธฐ, ์—†์œผ๋ฉด ๊ธฐ๋ณธ๊ฐ’์œผ๋กœ http://mlflow:5000 ์‚ฌ์šฉ
mlflow_uri = os.getenv("MLFLOW_TRACKING_URI", "http://mlflow:5000") # deployment์—์„œ mlflow tracking uri ๊ฐ€์ ธ์˜ค๊ธฐ, ์—†์œผ๋ฉด ๊ธฐ๋ณธ๊ฐ’์œผ๋กœ http://mlflow:5000 ์‚ฌ์šฉ
mlflow.set_tracking_uri(mlflow_uri)

model_uri = f"models:/{model_name}/latest"
model = mlflow.sklearn.load_model(model_uri)

_model_cache[cache_key] = model
return model


# mlflow์—์„œ uplift ๋ชจ๋ธ ๋กœ๋“œํ•˜๋Š” ํ•จ์ˆ˜
def load_uplift_model(model_name: str):
cache_key = f"{model_name}_uplift"
if cache_key in _model_cache:
return _model_cache[cache_key]

mlflow_uri = os.getenv("MLFLOW_TRACKING_URI", "http://mlflow:5000")
mlflow.set_tracking_uri(mlflow_uri)

model_uri = f"models:/{model_name}/latest"
Expand Down
31 changes: 31 additions & 0 deletions services/model-api/app/models/uplift_model.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
import numpy as np
from sklearn.base import BaseEstimator
from sklearn.linear_model import LogisticRegression
from sklearn.pipeline import Pipeline
from sklearn.preprocessing import StandardScaler

# ๊ฐ ๊ทธ๋ฃน์„ ๋กœ์ง€์Šคํ‹ฑ ํšŒ๊ท€๋ชจ๋ธ๋กœ ํ•™์Šต ํ›„ ์˜ˆ์ธก๊ฐ’์˜ ์ฐจ์ด๋กœ uplift ์ ์ˆ˜ ๊ณ„์‚ฐํ•˜๋Š” T-learner ๋ชจ๋ธ
class TLearner(BaseEstimator):
def __init__(self, C=1.0):
self.C = C
# treatment(์‚ฌ์šฉ์ž์—๊ฒŒ ํ”„๋กœ๋ชจ์…˜์„ ์ œ๊ณตํ•œ ๊ทธ๋ฃน)๊ณผ control(ํ”„๋กœ๋ชจ์…˜์„ ์ œ๊ณตํ•˜์ง€ ์•Š์€ ๊ทธ๋ฃน) ๊ฐ๊ฐ์— ๋Œ€ํ•ด ๋ชจ๋ธ์„ ํ•™์Šต
self.treatment_model = Pipeline([
('scaler', StandardScaler()),
('clf', LogisticRegression(C=C, random_state=42, max_iter=1000)),
])
self.control_model = Pipeline([
('scaler', StandardScaler()),
('clf', LogisticRegression(C=C, random_state=42, max_iter=1000)),
])
# fit ๋ฉ”์„œ๋“œ์—์„œ๋Š” treatment ๊ทธ๋ฃน๊ณผ control ๊ทธ๋ฃน์„ ๋‚˜๋ˆ„์–ด ๊ฐ๊ฐ์˜ ๋ชจ๋ธ์„ ํ•™์Šต
def fit(self, X, y, treatment):
treatment = np.array(treatment)
self.treatment_model.fit(X[treatment == 1], y[treatment == 1])
self.control_model.fit(X[treatment == 0], y[treatment == 0])
return self
# predict ๋ฉ”์„œ๋“œ์—์„œ๋Š” ๋‘ ๋ชจ๋ธ์˜ ์˜ˆ์ธก๊ฐ’์˜ ์ฐจ์ด๋กœ uplift ์ ์ˆ˜๋ฅผ ๊ณ„์‚ฐํ•˜์—ฌ ๋ฐ˜ํ™˜
# (๋งˆ์ผ€ํŒ…์„ ํ•˜๋ฉด ๊ตฌ๋งคํ•  ํ™•๋ฅ ) - (๋งˆ์ผ€ํŒ…์„ ์•ˆ ํ•˜๋ฉด ๊ตฌ๋งคํ•  ํ™•๋ฅ ) = uplift ์ ์ˆ˜
def predict(self, X):
p_treatment = self.treatment_model.predict_proba(X)[:, 1]
p_control = self.control_model.predict_proba(X)[:, 1]
return p_treatment - p_control
3 changes: 2 additions & 1 deletion services/model-api/app/services/service_factory.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,10 @@
from app.services import mnist_service, churn_service
from app.services import mnist_service, churn_service, uplift_service


_services = {
"mnist": mnist_service,
"churn": churn_service,
"uplift": uplift_service,
}

# ๋ชจ๋ธ์˜ ์ข…๋ฅ˜์— ๋งž๊ฒŒ service ํŒŒ์ผ ๋ถˆ๋Ÿฌ์˜ค๋Š” ๋ผ์šฐํ„ฐ
Expand Down
43 changes: 43 additions & 0 deletions services/model-api/app/services/uplift_service.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
import time
from datetime import datetime, timezone

from app.models.loader import load_uplift_model

MODEL_NAME = "uplift"

FEATURES = [
'account_age_months',
'avg_order_value',
'total_orders',
'days_since_last_purchase',
'discount_usage_rate',
'return_rate',
'browsing_frequency_per_week',
'cart_abandonment_rate',
]


def predict(data: dict) -> dict:
missing = [f for f in FEATURES if f not in data]
if missing:
raise ValueError(f"๋ˆ„๋ฝ๋œ ํ”ผ์ฒ˜: {missing}")

model = load_uplift_model(MODEL_NAME) # ๋ชจ๋ธ ๋กœ๋“œ

x = [[data[f] for f in FEATURES]]

start = time.time()
uplift_score = float(model.predict(x)[0])
duration = time.time() - start

return {
"result": {
"uplift_score": uplift_score,
"is_persuadable": uplift_score > 0, # ๋งˆ์ผ€ํŒ… ์•ก์…˜์ด ํšจ๊ณผ๊ฐ€ ์žˆ๋Š” ๊ณ ๊ฐ์ธ์ง€ ์—ฌ๋ถ€
},
"metadata": {
"model": MODEL_NAME,
"inference_time_ms": round(duration * 1000, 3),
"timestamp": datetime.now(timezone.utc).isoformat(),
},
}
Loading
Loading