요약
Apache Airflow와 Spark를 연동하면 대규모 배치 데이터 처리를 자동화하고 모니터링할 수 있습니다. SparkSubmitOperator를 통해 Spark 애플리케이션을 DAG Task로 실행하고, Web UI에서 작업 상태를 추적할 수 있습니다. Docker Compose 환경에서 Spark 클러스터를 구성하고 Airflow와 통합하는 방법을 다룹니다.
Apache Spark
대규모 데이터를 빠르고 효율적으로 처리하는 분산 데이터 처리 프레임워크
주요 특징
- In-Memory Computing
- 메모리(RAM)에서 처리하여 디스크 I/O가 많은 Hadoop보다 훨씬 빠름
- 메모리에 유지한 채 연산 가능 → 머신러닝, 데이터 변환에 최적
- 다양한 데이터 처리 방식 지원
- RDD (기본적인 데이터 구조)
- DataFrame
- Spark SQL
- Scalability (확장성)
- 수십~수천 대 클러스터 노드에서 병렬 실행 가능
- AWS, Azure, Google Cloud에서도 확장 가능
- 배치 & 실시간 데이터 처리 모두 가능
- 기본은 배치 처리지만 스트리밍 처리도 지원
Apache Spark 배치 처리 방식
RDD (Resilient Distributed Dataset)
- 분산된 데이터를 저장하고 처리하는 기본 단위
- Immutable(불변성) 지원 → 안정적 처리
- 여러 노드에서 병렬 처리 가능
DataFrame
- 구조화된 데이터를 최적화하여 처리
- Pandas와 유사
- Spark SQL과 연동 → SQL 쿼리로 데이터 변환 가능
Spark SQL
- SQL을 통해 데이터를 쉽게 조회하고 변환
- 다양한 데이터 소스 (HDFS, S3, JDBC, Hive, Cassandra 등) 연결 가능
Apache Spark의 RDD 구조 및 처리 과정
- 여러 Worker Node에 분산 저장
- Driver Node가 RDD의 연산을 스케줄링하고 Worker Node에 작업 분배

처리 단계
- 데이터 소스를 parallelize하거나 외부 소스에서 읽어 RDD 생성
- Transformation 연산(map, filter 등)은 RDD를 새로 생성하지만 즉시 실행되진 않음
- 여러 Transformation이 체이닝되어 실행 계획(DAG)을 구성
- Action 연산(count, collect 등) 호출 시 DAG가 실행되어 실제 처리 수행

Apache Spark의 DataFrame 구조 및 처리 과정
- 다양한 소스(Hive, CSV, JSON, RDBMS, XML, Cassandra 등) 통합 → 구조화된 형식으로 표현
- Spark SQL을 통해 생성된 DataFrame은 열(Column) 기반 테이블 형태

처리 단계
- 다양한 데이터 소스에서
read()또는load()로 DataFrame 생성 - 생성된 DataFrame에
select(),filter(),groupBy()등의 Transformation 적용 - Transformation이 체이닝되어 새 DataFrame을 연속 생성
- 마지막에
show(),count(),write()같은 Action 호출로 실제 실행

Apache Spark 배치 처리 주요 활용 사례
| 구분 | 활용 사례 설명 |
|---|---|
| 데이터 웨어하우스 적재 (ETL) | 매일 수집한 CSV 데이터를 Parquet 변환 후 Snowflake, BigQuery 적재 |
| 로그 데이터 분석 | 웹/앱/서버 로그를 분석하여 사용자 행동 분석 (예: 클릭 로그 분석으로 마케팅 전략 최적화) |
| 머신러닝 데이터 전처리 | 대량 데이터를 전처리 후 머신러닝 모델 학습에 활용 (예: 추천 시스템용 사용자 행동 데이터 전처리) |
Airflow에서 Spark를 활용하는 이유

핵심 이유
- 단순 Python 코드로는 처리 어려운 대규모 데이터 처리
- Pandas 한계 존재 (메모리 로딩, 병렬 처리 어려움)
- Spark를 활용한 대량 데이터 ETL 수행
- 복잡한 ETL 작업을 효율적으로 처리
- Spark 작업 모니터링 및 실패 시 자동 재시도
- Airflow는 Spark 작업 실패를 감지하고 자동 재시도 가능
- Airflow Web UI를 통해 Task 로그 확인, 오류 파악 가능
배치 워크플로우에서 DAG의 중요성
자동화
- 사람이 직접 실행할 필요 없이 정해진 스케줄에 따라 자동 실행
유지보수 용이
- 실행 로그 및 실패 이력을 기록하여 문제 해결 가능
확장성
- 여러 작업(Task)을 DAG 안에서 관리
- 병렬 실행을 통해 전체 워크플로우 시간 단축
재시도 및 오류 감지
- Task 실패 시 자동 재시도
- 안정적이고 신뢰성 있는 워크플로우 운영 가능
배치 워크플로우에서 DAG 처리 흐름 예시
DAG 구성
- 병렬 Task와 순차 Task 조합 가능
- 데이터 로드 → 사전 처리 → 후속 처리 흐름
병렬 처리
- 작업 간 의존성을 명확히 설정하여 병렬로 실행
- 전체 처리 시간을 단축할 수 있음

DAG에서 CSV 파일 처리 작업 실행
PythonOperator 활용
- PythonOperator를 사용하여 CSV 파일 읽기 & 변환
- 예시 작업
input_data.csv를 읽음- 컬럼명 변경 (old_column → new_column)
process_date컬럼 추가- 결과를
output_data.csv로 저장
- Airflow Web UI에서 작업 성공 여부와 로그 확인 가능
BashOperator를 활용한 데이터 전처리 스크립트 실행
- Shell 스크립트를 BashOperator로 실행
- 예시 작업
- 데이터 파일 정렬
- 중복 제거
- 데이터 포맷 수정
- Airflow Web UI에서 작업 진행 상황 및 로그를 모니터링 가능
병렬 실행 및 성능 최적화 방법
- 병렬 실행(Parallel Execution)
- DAG 설계 시 병렬로 처리할 수 있는 Task들을 병렬 배치
- 성능 최적화
- 병렬 Task 설정을 통해 전체 작업 시간을 단축
- DAG 그래프 뷰에서 병렬 작업 흐름 시각적으로 확인 가능
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 증가
다음 단계
- Introduction - Airflow 기본 개념 및 설치 복습
- DAG & Task - Task 의존성 및 스케줄링 심화
- Advanced Skill - Dynamic DAG로 여러 Spark 작업 자동 생성
- Spark - Spark 심화 학습
- Batch vs Real-time - 배치/실시간 처리 개념