상위: Spark
요약
Apache Airflow와 Spark를 연동하면 대규모 배치 데이터 처리를 자동화하고 모니터링할 수 있습니다. SparkSubmitOperator를 통해 Spark 애플리케이션을 DAG Task로 실행하고, Web UI에서 작업 상태를 추적할 수 있습니다.
Airflow에서 Spark를 활용하는 이유
핵심 이유
- 단순 Python 코드로는 처리 어려운 대규모 데이터 처리
- Pandas 한계 존재 (메모리 로딩, 병렬 처리 어려움)
- Spark를 활용한 대량 데이터 ETL 수행
- 복잡한 ETL 작업을 효율적으로 처리
- Spark 작업 모니터링 및 실패 시 자동 재시도
- Airflow는 Spark 작업 실패를 감지하고 자동 재시도 가능
- Airflow Web UI를 통해 Task 로그 확인, 오류 파악 가능
배치 워크플로우에서 DAG의 중요성
자동화
- 사람이 직접 실행할 필요 없이 정해진 스케줄에 따라 자동 실행
유지보수 용이
- 실행 로그 및 실패 이력을 기록하여 문제 해결 가능
확장성
- 여러 작업(Task)을 DAG 안에서 관리
- 병렬 실행을 통해 전체 워크플로우 시간 단축
재시도 및 오류 감지
- Task 실패 시 자동 재시도
- 안정적이고 신뢰성 있는 워크플로우 운영 가능
SparkSubmitOperator
SparkSubmitOperator 개요 및 역할
- Apache Airflow에서 Spark 애플리케이션을 실행하기 위한 전용 연산자
- Spark 클러스터(YARN, Kubernetes 등)에
.py,.jar,.scala파일 등을 제출(submit) - 복잡한 Spark 작업을 Airflow DAG 안에 통합하여 자동화된 데이터 파이프라인 구성 가능
- DAG 태스크로 Spark 작업 실행을 스케줄링 및 추적 가능
- SparkSubmit 명령어를 Python 코드로 대체하여 운영 효율성 확보
SparkSubmitOperator 주요 파라미터
| 파라미터 | 설명 |
|---|---|
| application | 실행할 Spark 애플리케이션 경로 (.py, .jar 등) |
| conf | Spark 설정 (예: "spark.executor.memory": "2g") |
| executor_memory, driver_memory | 리소스 지정 |
| application_args | 애플리케이션에 전달할 인자 목록 |
| conn_id | Spark 클러스터 연결 정보 (예: spark_default) |
| name | Spark 작업 이름 |
Environment Setting
docker-compose.yaml에 Spark 추가
spark-master:
image: bitnami/spark:latest
container_name: spark-master
environment:
- SPARK_MODE=master
- SPARK_RPC_AUTHENTICATION_ENABLED=no
- SPARK_RPC_ENCRYPTION_ENABLED=no
- SPARK_LOCAL_STORAGE_ENCRYPTION_ENABLED=no
- SPARK_SSL_ENABLED=no
ports:
- "8081:8080"
- "7077:7077"
networks:
- airflow
volumes:
- ${AIRFLOW_PROJ_DIR:-.}/airflow/dags/scripts:/opt/airflow/dags/scripts
spark-worker:
image: bitnami/spark:latest
container_name: spark-worker
environment:
- SPARK_MODE=worker
- SPARK_MASTER_URL=spark://spark-master:7077
depends_on:
- spark-master
ports:
- "8082:8081"
networks:
- airflow
volumes:
- ${AIRFLOW_PROJ_DIR:-.}/airflow/dags/scripts:/opt/airflow/dags/scripts
SparkSubmitOperator를 위한 Spark Provider 설치
airflow-webserver:
...
environment:
- _PIP_ADDITIONAL_REQUIREMENTS=apache-airflow-providers-apache-spark
- JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64
airflow-scheduler:
...
environment:
- _PIP_ADDITIONAL_REQUIREMENTS=apache-airflow-providers-apache-spark
- JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64
airflow-worker:
...
environment:
- _PIP_ADDITIONAL_REQUIREMENTS=apache-airflow-providers-apache-spark
- JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64
airflow-triggerer:
...
environment:
- _PIP_ADDITIONAL_REQUIREMENTS=apache-airflow-providers-apache-spark
- JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64
JAVA_HOME 세팅 및 경로 설정
JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64
- Worker, Webserver, Scheduler, Triggerer 모두 위와 같이 환경 변수 설정.
Docker 네트워크 설정
networks:
- airflow
- spark-master, spark-worker, airflow-webserver 등 모든 서비스는 airflow 네트워크에 연결되어야 함.
DAG/scripts 경로 & volume 마운트 설정
volumes:
- ${AIRFLOW_PROD_DIR:-.}/airflow/dags/scripts:/opt/airflow/dags/scripts
/airflow/dags/scripts/경로에 spark 스크립트(spark_wordcount.py) 존재.
$ ls -al ~/airflow/dags/scripts
-rw-r--r-- 1 user user 639 Mar 28 19:45 spark_wordcount.py
Java 수동 설치 (Worker, Webserver, Scheduler 컨테이너)
$ sudo docker ps
$ sudo docker exec -u root -it <컨테이너명> bash
# apt-get update && apt-get install -y openjdk-17-jdk
- Webserver, Worker, Scheduler 컨테이너 각각 진입해서 설치해야 함.
Airflow UI에서 Spark Connection 추가
- Airflow UI → Admin → Connections
- 새 Connection 추가 (+ 버튼)
- 설정 값:
| 필드 | 값 |
|---|---|
| Connection Id | spark_default |
| Connection Type | Spark |
| Host | spark://spark-master |
| Port | 7077 |
| Deploy mode | client |
| Spark binary | spark-submit |
DAG 파일 - spark_submit_example.py
from airflow import DAG
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from datetime import datetime
with DAG(
dag_id='spark_submit_example',
start_date=datetime(2024, 1, 1),
schedule_interval=None,
catchup=False,
tags=['spark'],
) as dag:
submit_job = SparkSubmitOperator(
task_id='spark_submit_task',
application='/opt/airflow/dags/scripts/spark_wordcount.py',
conn_id='spark_default',
conf={'spark.master': 'spark://spark-master:7077'},
verbose=True
)
Spark 작업 파일 - spark_wordcount.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, avg
spark = SparkSession.builder \
.appName("DataFrameTest") \
.master("spark://spark-master:7077") \
.getOrCreate()
# 샘플 데이터
data = [
("Alice", "Math", 85),
("Bob", "Math", 90),
("Alice", "English", 78),
("Bob", "English", 83),
("Alice", "Science", 92),
("Bob", "Science", 87)
]
columns = ["name", "subject", "score"]
# DataFrame 생성
df = spark.createDataFrame(data, columns)
# 평균 점수 계산
avg_scores = df.groupBy("name").agg(avg("score").alias("average_score"))
# 결과 출력
avg_scores.show()
spark.stop()
Log based Debugging
Airflow Web UI 로그 확인 방법
- DAG 실행 중 Task 실패 시 Logs 탭에서 상세 에러 확인
- 로그에서 오류 원인 분석 (Python Exception, ConnectionError 등)
DAG 실행 로그 저장 위치
-
Airflow의
logs/폴더에 DAG 별로 로그가 저장됨 -
예시:
ls -al ~/airflow/logs
특정 DAG/Task 실행 로그 확인
-
로그 파일 이름 규칙
~/airflow/logs/<dag_id>/<task_id>/<execution_date>/...
에러 예시
ConnectionError발생store_dataTask 실패- 로그를 통해 실패 위치와 에러 메세지를 확인하고 디버깅
트러블슈팅
자주 발생하는 문제
Java 관련 에러
JAVA_HOME is not set
해결책: 모든 Airflow 컨테이너에 Java 설치 및 JAVA_HOME 환경변수 설정
Connection 에러
Connection 'spark_default' not found
해결책: Airflow UI에서 Spark Connection 추가 확인
네트워크 에러
Cannot connect to spark://spark-master:7077
해결책: docker-compose.yaml에서 네트워크 설정 확인, 모든 서비스가 같은 네트워크에 속해야 함
메모리 부족
java.lang.OutOfMemoryError
해결책: Spark 설정에서 executor_memory, driver_memory 증가
다음 단계
- Airflow - Airflow 기본 개념 복습
- DAG - Task 의존성 및 스케줄링 심화
- Batch vs Real-time - 배치/실시간 처리 개념