전체 그래프
Tools

Airflow & Spark

data-engineeringtool/airflowtool/spark

상위: Airflow · 관련: Spark

요약

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에 작업 분배

Airflow & Spark

처리 단계

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

Airflow & Spark

Apache Spark의 DataFrame 구조 및 처리 과정

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

Airflow & Spark

처리 단계

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

Airflow & Spark

Apache Spark 배치 처리 주요 활용 사례

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

Airflow에서 Spark를 활용하는 이유

Airflow & Spark

핵심 이유

  • 단순 Python 코드로는 처리 어려운 대규모 데이터 처리
    • Pandas 한계 존재 (메모리 로딩, 병렬 처리 어려움)
  • Spark를 활용한 대량 데이터 ETL 수행
    • 복잡한 ETL 작업을 효율적으로 처리
  • Spark 작업 모니터링 및 실패 시 자동 재시도
    • Airflow는 Spark 작업 실패를 감지하고 자동 재시도 가능
    • Airflow Web UI를 통해 Task 로그 확인, 오류 파악 가능

배치 워크플로우에서 DAG의 중요성

자동화

  • 사람이 직접 실행할 필요 없이 정해진 스케줄에 따라 자동 실행

유지보수 용이

  • 실행 로그 및 실패 이력을 기록하여 문제 해결 가능

확장성

  • 여러 작업(Task)을 DAG 안에서 관리
  • 병렬 실행을 통해 전체 워크플로우 시간 단축

재시도 및 오류 감지

  • Task 실패 시 자동 재시도
  • 안정적이고 신뢰성 있는 워크플로우 운영 가능

배치 워크플로우에서 DAG 처리 흐름 예시

DAG 구성

  • 병렬 Task와 순차 Task 조합 가능
  • 데이터 로드 → 사전 처리 → 후속 처리 흐름

병렬 처리

  • 작업 간 의존성을 명확히 설정하여 병렬로 실행
  • 전체 처리 시간을 단축할 수 있음

Airflow & Spark

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

다음 단계

  • Introduction - Airflow 기본 개념 및 설치 복습
  • DAG & Task - Task 의존성 및 스케줄링 심화
  • Advanced Skill - Dynamic DAG로 여러 Spark 작업 자동 생성
  • Spark - Spark 심화 학습
  • Batch vs Real-time - 배치/실시간 처리 개념

참고 자료