[SI] 프로젝트 7주차

소복치·2025년 2월 13일

이번주 목표는 내가 맡은 업무에 대한 테이블을 생성하고
아직 실데이터를 받지 못했기 때문에, 샘플 데이터를 만들어서
코드가 잘돌아가는지 확인하는 작업을 진행할 예정이다.

후에, Airflow라는 스케줄링 오픈소스를 이용하여 배치JOB을 해볼 것이다.

이번주에는 몸상태가 좋지 않아 제대로된 일을 진행을 못했다..

테이블 생성 방법과 코드는 전 주차때 이미 작성이 완료된 상황이므로
이번 주차에서 설명은 링크로 대체 한다.

테이블 생성 방법 : https://velog.io/@jm4422/SI-%ED%94%84%EB%A1%9C%EC%A0%9D%ED%8A%B8-6%EC%A3%BC%EC%B0%A8
코드 예시 : https://velog.io/@jm4422/SI%ED%94%84%EB%A1%9C%EC%A0%9D%ED%8A%B8-5%EC%A3%BC%EC%B0%A8

1. Airflow란

워크플로우 작성, 예약 및 모니터링을 프로그래밍 방식으로 할 수 있는 오픈소스 플랫폼
워크 플로우란? 의존성으로 연결된 작업(Task)들의 집합

2. Airflow 기본 구성 및 작동 원리

  • DAGs : 데이터 엔지니어가 작성한 워크플로우 정의 파일
  • UI : 웹 기반 대시보드, DAG 상태 및 로그 모니터링
  • Web Server : UI롸 연결된 HTTP 서버로, 사용자의 요청을 처리
  • Scheduler : DAG을 실행할 Task의 순서를 정하고 Exeutor에 전달
  • Executor : Task 실행을 담당
  • Worker(s) : 실제 Task를 수행하는 작업 노드
  • DB : DAG 실행 이력 및 상태를 저장하는 데이터베이스
  • Airflow.cfg : 전체 환경 및 실행 설정 파일

🚀 Airflow 실행 흐름

1️⃣ Data Engineer가 DAG 작성

2️⃣ Web Server를 통해 DAG 모니터링

3️⃣ Scheduler가 DAG 실행을 결정

4️⃣ Executor가 Task를 Worker로 전달

5️⃣ Worker가 실제 작업 수행

6️⃣ 사용자는 UI에서 DAG 상태 확인

3. DAG 작성 방법 예시

from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from datetime import datetime

# 기본 설정
default_args = {
    'owner': 'data_engineer',
    'start_date': datetime(2024, 2, 14),
    'retries': 1,
}

dag = DAG(
    'example_dag',
    default_args=default_args,
    schedule_interval='@daily'
)

# 실행할 함수 정의
def print_message():
    print("Hello, Airflow!")

# Task 정의
task = PythonOperator(
    task_id='내가 스케줄링을 돌릴 id명',
    python_callable=print_message,
    dag=dag
)
  • Scheduler가 DAG을 감지하고 Executor에 Task 실행 요청
  • Worker가 print_message() 함수를 실행
  • 실행 결과가 Metadata DB에 저장
  • Web Server에서 실행 상태 확인 가능

간단하게 다시 흐름을 정리한다면,
DAG 작성 → Scheduler 실행 → Executor가 Task 배포 → Worker가 실행 → UI에서 모니터링

이렇게 진행이 된다.
airflow로 스케줄링 관리를 한다면, 내가 원하는 시간에 데이터 적재가 가능해지고
현재 데이터가 잘 들어가는지 로그 확인도 UI화면에서 가능하기 때문에 매우 편리하다.

이처럼 말로만 들어봤던 Airflow를 실전에서 한번 써볼 수 있는 기회가 있게되어서 너무 좋았다.

⬇️⬇️⬇️⬇️ Airflow UI 화면 예시 ⬇️⬇️⬇️⬇️

profile
오늘 터져도내일 다시극복

0개의 댓글