from databricks.sdk import WorkspaceClient
from databricks.sdk.service.jobs import (
Task, NotebookTask, TaskDependency,
CronSchedule, PauseStatus,
)
w = WorkspaceClient()
job_name = "LGIT_MLOps_Level2_AutoRetrain_Pipeline"
tasks = [
Task(task_key="feature_engineering",
notebook_task=NotebookTask(notebook_path=f"{notebook_base}/02_structured_feature_engineering")),
Task(task_key="model_training",
depends_on=[TaskDependency(task_key="feature_engineering")],
notebook_task=NotebookTask(notebook_path=f"{notebook_base}/03_structured_model_training")),
Task(task_key="model_registration",
depends_on=[TaskDependency(task_key="model_training")],
notebook_task=NotebookTask(notebook_path=f"{notebook_base}/04_model_registration_uc")),
Task(task_key="challenger_validation",
depends_on=[TaskDependency(task_key="model_registration")],
notebook_task=NotebookTask(notebook_path=f"{notebook_base}/05_challenger_validation")),
Task(task_key="batch_inference",
depends_on=[TaskDependency(task_key="challenger_validation")],
notebook_task=NotebookTask(notebook_path=f"{notebook_base}/06_batch_inference")),
Task(task_key="model_monitoring",
depends_on=[TaskDependency(task_key="batch_inference")],
notebook_task=NotebookTask(notebook_path=f"{notebook_base}/08_model_monitoring")),
Task(task_key="auto_retrain_if_drift",
depends_on=[TaskDependency(task_key="model_monitoring")],
notebook_task=NotebookTask(notebook_path=f"{notebook_base}/03d_retraining_strategies")),
]
created_job = w.jobs.create(
name=job_name,
tasks=tasks,
schedule=CronSchedule(
quartz_cron_expression="0 0 2 ? * MON", # 매주 월요일 02:00 KST
timezone_id="Asia/Seoul",
pause_status=PauseStatus.PAUSED,
),
tags={"project": "lgit-mlops-poc", "level": "2", "type": "auto-retrain"},
max_concurrent_runs=1,
)