CDE 의 CDEJobRunOperator 는 Spark Job 을 별도 Pod 에서 실행하므로 DAG 코드에서 JDBC 를 직접 호출할 수 없다. Airflow Connection(mariadb_jdbc)에 저장한 접속 정보를 PythonOperator 가 꺼내 XCom 으로 넘기고, CDEJobRunOperator 의 overrides.spark.args 에서 Jinja 로 받아 Spark Job 의 sys.argv 로 전달한다.
CDE Editor 에서 Operator 두 개로 구성하면 DAG 파일을 따로 만들 필요가 없다. PythonOperator 의 inline code 는 다음과 같다.
from airflow.hooks.base import BaseHook
from airflow.operators.python import get_current_context
CONN_ID = "mariadb_jdbc"
QUERY = "SELECT * FROM your_table LIMIT 10"
context = get_current_context()
conn = BaseHook.get_connection(CONN_ID)
context['ti'].xcom_push(key='jdbc', value={
'host': conn.host,
'port': str(conn.port or 3306),
'database': conn.schema,
'user': conn.login,
'password': conn.password,
'query': QUERY,
})
CDEJobRunOperator(task_id spark_jdbc_clt, Job Name 은 등록한 Spark Job)의 args 에는 {{ ti.xcom_pull(task_ids="python_2", key="jdbc")["host"] }} 형태로 항목을 나열한다. inline code 안에서 context 변수가 자동 주입되는지는 확인되지 않았으므로 get_current_context() 를 쓰는 편이 안전하다. 함수로 감싸서 실행되는 환경이면 return 값이 return_value 키로 push 되어 {{ (ti.xcom_pull(task_ids="python_2") | from_json)["host"] }} 로 읽을 수도 있다.
Spark Job 쪽은 argparse 로 6 개 인자를 받고 dbtable 옵션에 (query) AS t 서브쿼리를 넣는다. 등록 순서는 Resource 에 spark_jdbc_test.py 와 JDBC JAR 업로드 → Spark Job 등록(Jars: file:///app/mount/jars/mysql-connector-j-8.0.33.jar) → Airflow Connection 등록 → Airflow Job 을 Editor 로 구성이다.
argparse 를 쓰면 안 된다. 스케줄러가 DAG 를 파싱할 때 sys.argv 는 Airflow 자신의 인자라 파서가 깨진다. 재사용 파라미터는 Airflow Variable 이나 inline code 상단 상수로 둔다.import 하면 CDE 의 DAG validation 시점에 Resource 가 아직 mount 되지 않아 ModuleNotFoundError 가 난다. PythonOperator 안에서 sys.path.insert(0, '/app/mount') 후 동적으로 import 한다.export CDE_PIPELINE=... 로 만든 이름 spark-$CDE_PIPELINE-01-... 이 DAG 파이썬 런타임까지 전달되지 않아 $ 가 그대로 들어가고, task_id 는 영숫자 · - · . · _ 만 허용하므로 실패한다. os.environ.get("CDE_PIPELINE", "save_user_info") 로 읽거나 하드코딩한다.execute_jdbc 에서 ClassNotFoundException 이 나면 코드가 아니라 JVM classpath 문제다. Job 을 만들 때 JAR 을 지정한다. spark.jars 는 세션 생성 시점에만 적용되므로 코드 안에서 spark.conf.set 으로 넣어도 늦다.
cde job create \
--name spark-save_user_info-02-export \
--type spark \
--application-file save_user_info/spark-save_user_info-02-export.py \
--mount-1-resource ${CDE_RESOURCE} \
--jar /app/mount/jars/mysql-connector-j-8.0.33.jar