CDE(Cloudera Data Engineering) 위에서 데이터 파이프라인을 여러 개 운영할 때의 프로젝트 구성과, 외부에서 넘긴 값이 Spark Job 까지 어떤 경로로 도달하는지 정리한다. 역할 분담은 단순하다 — 데이터 처리는 Spark Job 이 하고, 실행 순서 제어는 Airflow DAG 이 하며, DAG 실행은 외부에서 CDE CLI 로 촉발한다.
CDE 의 가상 클러스터 · Job · Resource 구분은 CDE 구성 모델 에 있다.
NiFi (수집)
│ bronze 적재 후
│ cde job run --name dag-<name> --config-json '{"part_ymd":"...", "base_ymd":"..."}'
▼
CDE Airflow DAG (dag-<name>)
│ CDEJobRunOperator 로 Spark Job 을 순서대로 실행
├─▶ spark-<name>-01-silver-<table> 전처리
├─▶ spark-<name>-02-gold-<table> 최종 가공
└─▶ spark-<name>-03-export-<table> 대상 DB 적재
레이어의 뜻은 다음과 같다.
| 레이어 | 담당 |
|---|---|
| bronze | 원천 수집·저장. 수집 도구(NiFi)가 처리하는 것이 원칙 |
| silver | 비즈니스 로직 적용, 전처리 |
| gold | 대상 DB 인터페이스용 최종 가공 |
| export | gold 를 외부 DB 에 적재. 필요한 파이프라인만 |
두 값이 파이프라인 전체를 관통한다.
| 파라미터 | 용도 |
|---|---|
BASE_YMD |
원천 DB 조회 조건(기준 날짜) |
PART_YMD |
수집일. Hive 파티션 키 |
외부 호출
│ cde job run --config-json '{"part_ymd":"20260308","base_ymd":"20260307"}'
▼
CDE REST API
│ dag_run.conf = {"part_ymd": "...", "base_ymd": "..."}
▼
Airflow DAG
│ overrides.spark.conf 에 Jinja 로 렌더링
│ "spark.cde.part_ymd": "{{ dag_run.conf.get('part_ymd', ds_nodash) }}"
│ "spark.cde.base_ymd": "{{ dag_run.conf.get('base_ymd', ds_nodash) }}"
▼
Spark Job
spark.conf.get("spark.cde.part_ymd")
dag_run.conf.get(key, ds_nodash) 형태로 기본값을 두면, 값 없이 수동 실행해도 해당 실행일 기준으로 돈다.
DB 접속 정보도 같은 경로를 탄다. Airflow Connection 에 등록해 두고 DAG 에서 꺼내 overrides.spark.conf 에 실어 보낸다. DAG 코드에서 JDBC 를 직접 호출하지 않는다 — Spark Job 은 별도 파드에서 돌기 때문이다. 이 제약과 XCom 을 쓰는 변형은 CDE Airflow DAG 작성 시 자주 걸리는 것 에 있다.
Airflow Connection (conn_id)
│ BaseHook.get_connection(conn_id)
▼
overrides.spark.conf:
"spark.cde.src.db.host": "<host>"
"spark.cde.src.db.port": "<port>"
"spark.cde.src.db.schema": "<schema>"
▼
Spark Job 공통 모듈이 읽어 컨텍스트 객체로 제공
접속 비밀번호를 평문으로 spark.conf 에 싣지 않는다. 암호화해 넘기는 방식은 CDE Secret 을 DAG 에서 Spark 로 넘기기 를 본다.
_t1 = CDEJobRunOperator(task_id="spark-sales-01-silver-daily",
job_name="spark-sales-01-silver-daily",
overrides=_overrides)
_t2 = CDEJobRunOperator(task_id="spark-sales-02-gold-daily",
job_name="spark-sales-02-gold-daily",
overrides=_overrides)
_t1 >> _t2
job_name 은 CDE 에 등록된 Job 이름이다. task_id 와 같게 두면 장애 때 대조가 쉬워진다.
관리 도구와 CDE 에 올라가는 코드를 분리하는 것이 핵심이다.
cde-project/
├── scripts/ # 서버 전용. CDE 에 업로드하지 않는다
│ ├── pipeline_config.py # 파이프라인 정의 (steps, conn_id, params)
│ ├── dag_generator.py # 설정을 읽어 DAG 파일 생성
│ ├── spark_template.py # Spark Job 초기 템플릿 생성
│ ├── resource_setup.sh # 리소스 폴더에 utils/jars 복사
│ ├── resource_upload.sh # CDE Resource 업로드
│ ├── utils/spark_common.py # Spark Job 공통 모듈
│ └── jars/ # JDBC 드라이버 원본
├── resources/ # CDE 에 배포되는 코드
│ └── <리소스명>/
│ ├── <파이프라인명>/
│ │ ├── dag-<파이프라인명>.py
│ │ └── spark-<...>.py
│ ├── utils/ # scripts/utils 복사본
│ └── jars/ # scripts/jars 복사본
└── tests/
Resource 하나에 여러 Job 을 붙이는 것이 표준이다. DAG 마다 Resource 를 1:1 로 만들면 업로드 비용이 커지고 공통 모듈을 공유하기 어렵다.
| 대상 | 형식 | 예 |
|---|---|---|
| 파이프라인명 | _ 구분 |
sales, sales_export |
| Spark Job · 파일명 | spark-<파이프라인>-<순서>-<layer>-<테이블> |
spark-sales-01-silver-daily |
| DAG Job 명 | dag-<파이프라인명> |
dag-sales |
| 마운트 경로 | Resource 는 정해진 마운트 경로 아래에 붙는다 | sys.path.append("/app/mount/utils") |
Job 이름에 순서와 레이어가 들어 있으면 Job Runs 화면에서 파이프라인 진행 상황이 이름만으로 읽힌다.
1. pipeline_config.py 에 파이프라인 정의 추가 (steps, conn_id)
2. Spark Job 템플릿 생성
3. 비즈니스 로직 작성
4. DAG 파일 생성
5. 리소스 폴더 초기화 (utils/jars 복사)
6. CDE Resource 업로드 — 공통 모듈과 파이프라인 코드를 나눠 올린다
7. CDE 에 Spark Job 과 DAG Job 등록
8. 실행
cde job run --name dag-<name> --config-json '{"part_ymd":"...","base_ymd":"..."}'
DAG 파일을 손으로 쓰지 않고 설정에서 생성하면 파이프라인이 늘어나도 overrides 구성과 Task 연결 방식이 한 곳에서 관리된다. CLI 인증 설정은 CDE CLI 인증과 API 토큰 에 있다.