diff --git a/.github/workflows/cd-production.yaml b/.github/workflows/cd-production.yaml index 1b19550..9110f03 100644 --- a/.github/workflows/cd-production.yaml +++ b/.github/workflows/cd-production.yaml @@ -7,6 +7,7 @@ on: paths: - 'services/model-api/**' - 'services/alert-bridge/**' + - 'k8s/**' jobs: deploy: @@ -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 diff --git a/.github/workflows/deploy-churn.yaml b/.github/workflows/deploy-churn.yaml new file mode 100644 index 0000000..ba15928 --- /dev/null +++ b/.github/workflows/deploy-churn.yaml @@ -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 diff --git a/.github/workflows/deploy-uplift.yaml b/.github/workflows/deploy-uplift.yaml new file mode 100644 index 0000000..e3ca279 --- /dev/null +++ b/.github/workflows/deploy-uplift.yaml @@ -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 diff --git a/.github/workflows/retrain-churn.yaml b/.github/workflows/retrain-churn.yaml deleted file mode 100644 index 944bcb0..0000000 --- a/.github/workflows/retrain-churn.yaml +++ /dev/null @@ -1,46 +0,0 @@ -name: Retrain (churn) - -on: - repository_dispatch: - types: - - model-retrain - -jobs: - retrain: - runs-on: ubuntu-latest - - steps: - - name: Checkout - uses: actions/checkout@v4 - - - name: Python 설치 - uses: actions/setup-python@v5 - with: - python-version: "3.11" - - - name: 의존성 설치 - run: | - pip install -r services/model-api/requirements.txt - - - name: 트리거된 알림 출력 - run: | - echo "Triggered by alerts: ${{ toJson(github.event.client_payload.alerts) }}" - - - name: 모델 재학습 & MLflow 등록 - working-directory: services/model-api - env: - MLFLOW_TRACKING_URI: https://mlflow.swmlops.site - AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }} - AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }} - AWS_DEFAULT_REGION: ap-northeast-2 - run: | - python training/churn/train.py \ - --mlflow_uri https://mlflow.swmlops.site \ - --data_dir ../data/churn - - - name: model-api 파드 재시작 - run: | - mkdir -p $HOME/.kube - echo "${{ secrets.KUBECONFIG_PRODUCTION }}" | base64 -d > $HOME/.kube/config - kubectl rollout restart deployment/model-api -n model-api - kubectl rollout status deployment/model-api -n model-api diff --git a/k8s/alert-bridge/deployment.yaml b/k8s/alert-bridge/deployment.yaml index d8aa97d..07fb293 100644 --- a/k8s/alert-bridge/deployment.yaml +++ b/k8s/alert-bridge/deployment.yaml @@ -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 diff --git a/services/alert-bridge/main.py b/services/alert-bridge/main.py index 570c987..eea1d9f 100644 --- a/services/alert-bridge/main.py +++ b/services/alert-bridge/main.py @@ -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} diff --git a/services/alert-bridge/requirements.txt b/services/alert-bridge/requirements.txt index d23d558..97dc7cd 100644 --- a/services/alert-bridge/requirements.txt +++ b/services/alert-bridge/requirements.txt @@ -1,3 +1,2 @@ fastapi uvicorn -httpx diff --git a/services/model-api/app/main.py b/services/model-api/app/main.py index efaa355..3b62c3d 100644 --- a/services/model-api/app/main.py +++ b/services/model-api/app/main.py @@ -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), + ) diff --git a/services/model-api/app/models/loader.py b/services/model-api/app/models/loader.py index 02e6174..c567c16 100644 --- a/services/model-api/app/models/loader.py +++ b/services/model-api/app/models/loader.py @@ -2,6 +2,7 @@ import mlflow.pytorch import mlflow.sklearn +from app.models.uplift_model import TLearner # noqa: F401 MLflow 역직렬화 시 필요 _model_cache: dict = {} @@ -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" @@ -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" diff --git a/services/model-api/app/models/uplift_model.py b/services/model-api/app/models/uplift_model.py new file mode 100644 index 0000000..958d383 --- /dev/null +++ b/services/model-api/app/models/uplift_model.py @@ -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 diff --git a/services/model-api/app/services/service_factory.py b/services/model-api/app/services/service_factory.py index 93ba777..438d748 100644 --- a/services/model-api/app/services/service_factory.py +++ b/services/model-api/app/services/service_factory.py @@ -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 파일 불러오는 라우터 diff --git a/services/model-api/app/services/uplift_service.py b/services/model-api/app/services/uplift_service.py new file mode 100644 index 0000000..8bf48b1 --- /dev/null +++ b/services/model-api/app/services/uplift_service.py @@ -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(), + }, + } diff --git a/services/model-api/training/churn/train.py b/services/model-api/training/churn/train.py index dfd190e..097737e 100644 --- a/services/model-api/training/churn/train.py +++ b/services/model-api/training/churn/train.py @@ -7,6 +7,7 @@ import argparse import os import sys +import requests sys.path.insert(0, os.path.join(os.path.dirname(__file__), '../..')) @@ -58,6 +59,24 @@ def main(args): print(f"test_roc_auc: {test_metrics['roc_auc']:.4f}") print(f'MLflow에 모델 등록 완료: churn-{args.version}') + github_token = os.getenv("GITHUB_TOKEN") + github_owner = os.getenv("GITHUB_REPO_OWNER") + github_repo = os.getenv("GITHUB_REPO_NAME") + + if github_token and github_owner and github_repo: + response = requests.post( + f"https://api.github.com/repos/{github_owner}/{github_repo}/dispatches", + headers={ + "Authorization": f"Bearer {github_token}", + "Accept": "application/vnd.github+json", + }, + json={"event_type": "deploy-churn"}, + ) + if response.status_code == 204: + print("배포 트리거 완료") + else: + print(f"배포 트리거 실패: {response.status_code}") + if __name__ == '__main__': parser = argparse.ArgumentParser() diff --git a/services/model-api/training/uplift/dataset.py b/services/model-api/training/uplift/dataset.py new file mode 100644 index 0000000..e305d97 --- /dev/null +++ b/services/model-api/training/uplift/dataset.py @@ -0,0 +1,43 @@ +import numpy as np +import pandas as pd + +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 generate_dummy_data(n_samples=5000, random_state=42): + np.random.seed(random_state) + + df = pd.DataFrame({ + 'account_age_months': np.random.randint(1, 60, n_samples), + 'avg_order_value': np.random.uniform(10000, 200000, n_samples), + 'total_orders': np.random.randint(1, 50, n_samples), + 'days_since_last_purchase': np.random.randint(1, 365, n_samples), + 'discount_usage_rate': np.random.uniform(0, 1, n_samples), + 'return_rate': np.random.uniform(0, 0.5, n_samples), + 'browsing_frequency_per_week': np.random.uniform(0, 20, n_samples), + 'cart_abandonment_rate': np.random.uniform(0, 1, n_samples), + }) + + df['treatment'] = np.random.randint(0, 2, n_samples) + + # 이탈 위험이 높을수록 재방문 확률 낮음 + treatment 그룹에 uplift 효과 추가 + base_prob = ( + 0.3 + + 0.2 * (df['total_orders'] / 50) + - 0.2 * (df['days_since_last_purchase'] / 365) + - 0.1 * df['cart_abandonment_rate'] + ) + treatment_effect = 0.15 * df['treatment'] + prob = np.clip(base_prob + treatment_effect, 0.05, 0.95) + df['outcome'] = np.random.binomial(1, prob) + + return df diff --git a/services/model-api/training/uplift/evaluate.py b/services/model-api/training/uplift/evaluate.py new file mode 100644 index 0000000..0633cc3 --- /dev/null +++ b/services/model-api/training/uplift/evaluate.py @@ -0,0 +1,22 @@ +import numpy as np + +# 모델 평가 함수 +def evaluate(model, X, y, treatment): + treatment = np.array(treatment) + uplift_scores = model.predict(X) # uplift 점수 + + treatment_outcome_rate = y[treatment == 1].mean() # treatment 그룹의 재방문율 + control_outcome_rate = y[treatment == 0].mean() # control 그룹의 재방문율 + avg_uplift = uplift_scores.mean() # 전체 샘플의 평균 uplift 점수 + + high_uplift_mask = uplift_scores > np.percentile(uplift_scores, 75) # uplift 점수가 상위 25%인 샘플 + high_uplift_treatment_rate = y[(high_uplift_mask) & (treatment == 1)].mean() if (high_uplift_mask & (treatment == 1)).sum() > 0 else 0 # 실제 uplift 점수가 높은 샘플 중 treatment 그룹의 재방문율(uplift 점수 vs outcome) + high_uplift_control_rate = y[(high_uplift_mask) & (treatment == 0)].mean() if (high_uplift_mask & (treatment == 0)).sum() > 0 else 0 # 실제 uplift 점수가 높은 샘플 중 control 그룹의 재방문율(uplift 점수 vs outcome) + + return { + "treatment_outcome_rate": float(treatment_outcome_rate), + "control_outcome_rate": float(control_outcome_rate), + "avg_uplift_score": float(avg_uplift), + "high_uplift_treatment_rate": float(high_uplift_treatment_rate), + "high_uplift_control_rate": float(high_uplift_control_rate), + } diff --git a/services/model-api/training/uplift/train.py b/services/model-api/training/uplift/train.py new file mode 100644 index 0000000..4e8532d --- /dev/null +++ b/services/model-api/training/uplift/train.py @@ -0,0 +1,96 @@ +""" +Uplift 모델 학습 (T-Learner) + +# 학습 실행 +python3 training/uplift/train.py --version v1 +""" +import argparse +import os +import sys + +sys.path.insert(0, os.path.join(os.path.dirname(__file__), '../..')) + +import mlflow +import mlflow.sklearn +import numpy as np +import requests +from sklearn.model_selection import train_test_split + +from app.models.uplift_model import TLearner +from training.uplift.dataset import generate_dummy_data, FEATURES +from training.uplift.evaluate import evaluate + + +def main(args): + df = generate_dummy_data(n_samples=5000) # 일단 더미데이터로 학습, 실제로는 고객 데이터를 불러와서 사용 + + X = df[FEATURES].values + y = df['outcome'].values + treatment = df['treatment'].values + + X_train, X_test, y_train, y_test, t_train, t_test = train_test_split( + X, y, treatment, test_size=0.2, random_state=42 + ) + + mlflow.set_tracking_uri(args.mlflow_uri) + mlflow.set_experiment("uplift") + + with mlflow.start_run(run_name=args.version): + model = TLearner(C=args.C) + model.fit(X_train, y_train, t_train) # 모델 학습 + + metrics = evaluate(model, X_test, y_test, t_test) # 모델 평가 + + # mlflow 파라미터 저장 + mlflow.log_param("version", args.version) # 모델 버전 + mlflow.log_param("C", args.C) # 모델 하이퍼파라미터 + mlflow.log_param("n_samples", len(df)) # 학습에 사용된 데이터 샘플 수 + mlflow.log_param("features", FEATURES) # 모델에 사용된 피처 목록 + + # mlflow 메트릭 저장 + mlflow.log_metric("treatment_outcome_rate", metrics["treatment_outcome_rate"]) # treatment 그룹의 재방문율 + mlflow.log_metric("control_outcome_rate", metrics["control_outcome_rate"]) # control 그룹의 재방문율 + mlflow.log_metric("avg_uplift_score", metrics["avg_uplift_score"]) # 전체 샘플의 평균 uplift 점수 + mlflow.log_metric("high_uplift_treatment_rate", metrics["high_uplift_treatment_rate"]) # 실제 uplift 점수가 높은 샘플 중 treatment 그룹의 재방문율(uplift 점수 vs outcome) + mlflow.log_metric("high_uplift_control_rate", metrics["high_uplift_control_rate"]) # 실제 uplift 점수가 높은 샘플 중 control 그룹의 재방문율 + + # mlflow 모델 저장 및 등록 + mlflow.sklearn.log_model( + model, + artifact_path="model", + registered_model_name="uplift", + ) + + print(f"treatment_outcome_rate : {metrics['treatment_outcome_rate']:.4f}") + print(f"control_outcome_rate : {metrics['control_outcome_rate']:.4f}") + print(f"avg_uplift_score : {metrics['avg_uplift_score']:.4f}") + print(f"MLflow에 모델 등록 완료: uplift-{args.version}") + + # 배포 트리거 + github_token = os.getenv("GITHUB_TOKEN") + github_owner = os.getenv("GITHUB_REPO_OWNER") + github_repo = os.getenv("GITHUB_REPO_NAME") + + if github_token and github_owner and github_repo: + response = requests.post( + f"https://api.github.com/repos/{github_owner}/{github_repo}/dispatches", + headers={ + "Authorization": f"Bearer {github_token}", + "Accept": "application/vnd.github+json", + }, + json={"event_type": "deploy-uplift"}, + ) + if response.status_code == 204: + print("배포 트리거 완료") + else: + print(f"배포 트리거 실패: {response.status_code}") + + +if __name__ == '__main__': + parser = argparse.ArgumentParser() + parser.add_argument('--version', type=str, default='v1') + parser.add_argument('--C', type=float, default=1.0) + parser.add_argument('--mlflow_uri', type=str, default='https://mlflow.swmlops.site') + args = parser.parse_args() + + main(args)