전체 그래프
Spark

Spark Integration

data-engineeringairflowsparkintegration

상위: 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 등)
confSpark 설정 (예: "spark.executor.memory": "2g")
executor_memory, driver_memory리소스 지정
application_args애플리케이션에 전달할 인자 목록
conn_idSpark 클러스터 연결 정보 (예: spark_default)
nameSpark 작업 이름

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 추가

  1. Airflow UI → Admin → Connections
  2. 새 Connection 추가 (+ 버튼)
  3. 설정 값:
필드
Connection Idspark_default
Connection TypeSpark
Hostspark://spark-master
Port7077
Deploy modeclient
Spark binaryspark-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_data Task 실패
  • 로그를 통해 실패 위치와 에러 메세지를 확인하고 디버깅

트러블슈팅

자주 발생하는 문제

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 증가

다음 단계

참고 자료