FastAPI 가 만든 Parquet 파일을 S3 에 쌓고, Databricks Auto Loader 가 Bronze(Delta) 로 상시 적재하며, Airflow 가 주기적으로 Databricks Job 을 호출해 Bronze → Silver 변환을 수행하는 구성에서 DAG 를 일관되게 작성하기 위해 정한 규약이다. Claude 프로젝트 컨텍스트로 등록해 두고 코드 생성 시 전제로 쓰던 내용이다.
FastAPI ──Parquet──▶ S3 ──Auto Loader stream──▶ Bronze (Delta) ──Airflow cron──▶ Silver (Delta)
S3 → Bronze 는 Databricks Workflows 로 24/7 스트리밍하며 Airflow 가 개입하지 않는다. Bronze → Silver 만 Airflow cron 이 Databricks Job 을 트리거하고, 적재 방식(MERGE / OVERWRITE / APPEND)은 테이블마다 정한다. S3 파티셔닝은 데이터 특성에 따라 그때그때 정하며 고정 표준을 두지 않는다. 요청당 Parquet 파일 하나가 생기므로 small file 누적 문제가 있음을 전제로 운영한다.
extd_api_connection, Databricks 용 databricks_service_principal.| 항목 | 규칙 |
|---|---|
| 시간대 | pendulum + Asia/Seoul |
| DAG ID | <layer>_<domain> — 예: bronze_<db_name>, silver_<domain> |
| owner | 운영 계정 하나로 통일 (예: aip_ops) |
| 알림 | utils.notifications.teams.send_message 를 on_failure_callback / on_retry_callback 에 연결 |
| 시작·종료 | EmptyOperator start / end. end 는 trigger_rule="none_failed_min_one_success" |
| 병렬 처리 | 테이블 리스트를 fan-out — for table in tables: start >> pipeline >> end |
| 설정 분리 | dags/<layer>_pipeline/config/<layer>-<domain>-conf.yaml |
| 태그 | [<domain>, "databricks", "<layer>", "<team-tag>"] |
| TaskGroup | taskgroups.<카테고리>.<TaskGroupClass> — 예: ExtdToDatabricks |
| DAG 옵션 | max_active_runs=1, max_active_tasks=10, dagrun_timeout 명시 |
from datetime import timedelta
import pendulum
from airflow import DAG
from airflow.operators.empty import EmptyOperator
from taskgroups.<category> import <TaskGroupClass>
from utils.notifications.teams import send_message
db_name = ""
yaml_file_path = f"dags/<layer>_pipeline/config/<layer>-{db_name}-conf.yaml"
default_args = {
"owner": "aip_ops",
"retries": 1,
"retry_delay": timedelta(hours=1),
"on_failure_callback": send_message,
"on_retry_callback": send_message,
}
tables = ["", "", ""]
with DAG(
f"<layer>_{db_name}",
start_date=pendulum.datetime(2025, 10, 31, tz="Asia/Seoul"),
schedule="<cron>",
catchup=False,
default_args=default_args,
max_active_runs=1,
max_active_tasks=10,
dagrun_timeout=timedelta(hours=5),
tags=["<domain>", db_name, "databricks", "<layer>", "<team-tag>"],
doc_md=__doc__,
):
start = EmptyOperator(task_id="start")
end = EmptyOperator(task_id="end", trigger_rule="none_failed_min_one_success")
for table in tables:
pipeline = <TaskGroupClass>(
db_name, table,
<connection_kwargs>,
table_config_yaml_path=yaml_file_path,
data_check=True,
)
start >> pipeline >> end
테이블마다 TaskGroup 을 fan-out 하고 end 를 none_failed_min_one_success 로 두면, 한 테이블이 실패해도 나머지 테이블은 끝까지 진행되고 DAG 는 실패로 표시된다. max_active_runs=1 은 같은 DAG 의 이전 실행이 끝나기 전에 다음 스케줄이 겹쳐 돌아 MERGE 가 충돌하는 것을 막는다. 설정을 YAML 로 분리하면 테이블 추가·적재 방식 변경 때 DAG 코드를 건드리지 않아도 된다.
Authorization / Cookie / *token / *secret / *password 헤더는 마스킹한다.