# Lakeflow Jobs SDK로 위 파이프라인 생성
from databricks.sdk import WorkspaceClient
from databricks.sdk.service.jobs import *
w = WorkspaceClient()
w.jobs.create(
name="daily-data-pipeline",
tasks=[
Task(
task_key="ingest",
pipeline_task=PipelineTask(pipeline_id="<ingest-pipeline-id>"),
),
Task(
task_key="transform",
depends_on=[TaskDependency(task_key="ingest")],
pipeline_task=PipelineTask(pipeline_id="<sdp-pipeline-id>"),
),
Task(
task_key="quality_check",
depends_on=[TaskDependency(task_key="transform")],
sql_task=SqlTask(
query=SqlTaskQuery(query_id="<quality-check-query-id>"),
warehouse_id="<warehouse-id>"
),
),
Task(
task_key="update_features",
depends_on=[TaskDependency(task_key="quality_check")],
notebook_task=NotebookTask(
notebook_path="/Workspace/pipelines/update_features"
),
),
Task(
task_key="refresh_dashboard",
depends_on=[TaskDependency(task_key="quality_check")],
sql_task=SqlTask(
query=SqlTaskQuery(query_id="<refresh-mv-query-id>"),
warehouse_id="<warehouse-id>"
),
),
],
schedule=CronSchedule(
quartz_cron_expression="0 0 6 * * ?", # 매일 06:00
timezone_id="Asia/Seoul"
),
email_notifications=JobEmailNotifications(
on_failure=["data-team@company.com"]
)
)