[Airflow] Airflow 기반 초단기예보 수집 파이프라인 구축 기록

MJ·2026년 1월 19일

[1] 목표

  • 기상청 초단기 실황 API(getUltraSrtNcst)를 주기적으로 호출
  • Airflow DAG으로 스케줄링하여 시간 단위 데이터 자동 수집
  • 수집 결과를 JSON 파일로 저장하여 이후 Spark/분석 파이프라인과 연계

[2] 파이프라인 전체 흐름

  1. Airflow Scheduler가 매 정시 + 15분에 DAG 실행
  2. 실행 시점의 data_interval_end를 기준으로 조회 시간 계산
  3. Airflow Variable에 저장된 API Key를 사용해 기상청 API 호출
  4. 응답 결과를 JSON 파일로 로컬 디렉토리에 저장
  5. API 제약(과거 데이터 미제공)은 에러가 아닌 정상 스킵으로 처리

[3] 설계 시 고려사항

[3-1] 왜 Airflow를 사용했는가

  • 단순 크론이 아닌 데이터 파이프라인 단위 관리
  • 실행 이력, 재시도, 실패 로그를 UI에서 명확히 확인 가능
  • 이후 Spark, DB 적재, 후속 DAG과 자연스럽게 연결 가능

[3-2] 초단기예보 API의 제약 사항

  • 기상청 초단기 실황 API는 과거 일정 시점 이전 데이터 미제공
  • catchup 실행 시 과거 시간대 요청은 실패 가능
  • 따라서 “실패”가 아니라 “정상적으로 건너뜀”으로 처리해야 함

[4] DAG 핵심 설계 요소

[4-1] 시간 처리 전략

  • Airflow의 data_interval_end를 기준 시점으로 사용
  • 실제 API 호출 시각은 Asia/Seoul 타임존으로 변환
now = data_interval_end.in_timezone("Asia/Seoul")

이 방식의 장점:

  • 실행 시점과 데이터 기준 시점을 명확히 분리
  • catchup 실행 시에도 시간 계산 로직이 일관됨

[4-2] API Key 관리 방식

  • 코드에 하드코딩하지 않고 Airflow Variable 사용
api_key = Variable.get("weather_decoding_api_key")

이유:

  • Git 저장소에 민감 정보 노출 방지
  • 환경별(API Key 교체) 관리 용이
  • 운영 환경으로 확장 시에도 동일한 방식 유지 가능

[4-3] 데이터 저장 구조

  • Airflow 컨테이너의 include 디렉토리 사용
  • 시간 단위 파일명으로 저장
/usr/local/airflow/include/ultra_srt_weather/
└── ncst_YYYYMMDD_HH00.json

이유:

  • Airflow에서 공식적으로 권장되는 공유 경로
  • Spark, 후속 파이프라인에서 동일 경로 사용 가능

[5] DAG 코드 구조

[5-1] PythonOperator 기반 Task 구성

  • 단일 Task에서 API 호출 + 파일 저장 수행
  • 이후 Spark 처리 Task를 downstream으로 확장 가능
fetch_task = PythonOperator(
    task_id="fetch_ncst_step",
    python_callable=fetch_ultra_srt_ncst
)

[5-2] 스케줄링 설정

schedule="15 * * * *"
catchup=True
  • 매 정시 + 15분 실행
  • 장애 복구 시 과거 실행도 가능하도록 catchup 활성화
  • 과거 데이터 미제공 구간은 내부 로직에서 자동 스킵

[6] 전체 코드

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.decorators import dag, task

from airflow.models import Variable
from datetime import datetime, timedelta
import requests
import json
import os
import pendulum

local_tz = pendulum.timezone("Asia/Seoul")

def fetch_ultra_srt_ncst(data_interval_end, **kwargs):
    now = data_interval_end.in_timezone("Asia/Seoul")
    api_key = Variable.get("weather_decoding_api_key")
    url = 'http://apis.data.go.kr/1360000/VilageFcstInfoService_2.0/getUltraSrtNcst'
    save_dir = "/usr/local/airflow/include/ultra_srt_weather"
    os.makedirs(save_dir, exist_ok=True)

    params = {
        'serviceKey': api_key,
        'pageNo': '1', 'numOfRows': '1000', 'dataType': 'JSON',
        'base_date': now.strftime("%Y%m%d"),
        'base_time': now.strftime("%H00"), 
        'nx': '60', 'ny': '127'
    }

    response = requests.get(url, params=params, timeout=30)
    data = response.json()
    header = data['response']['header']

    if header['resultCode'] == '00':
        file_path = os.path.join(save_dir, f"ncst_{now.strftime('%Y%m%d_%H00')}.json")
        with open(file_path, 'w', encoding='utf-8') as f:
            json.dump(data, f, ensure_ascii=False, indent=4)
        print(f"저장 완료: {file_path}")
    else:
        # 1일이 지난 데이터는 기상청에서 주지 않으므로 에러로 처리하지 않고 그냥 넘어감
        print(f"건너뜀: {header['resultMsg']} ({now.strftime('%Y%m%d %H00')})")

with DAG(
    dag_id='kma_ultra_srt_ncst_collector_v4', # 새로운 ID
    default_args={'owner': 'user', 'retries': 1},
    # 현재 시간 기준 최근 12시간 전부터만 시도 (API 제약 고려)
    start_date=pendulum.now(local_tz).subtract(hours=12),
    schedule='15 * * * *', 
    catchup=True,
    tags=['kma', 'fixed']
) as dag:

    fetch_task = PythonOperator(
        task_id='fetch_ncst_step',
        python_callable=fetch_ultra_srt_ncst
    )

[7] 실제로 겪은 문제와 해결 과정

[7-1] 문제 1: Airflow Variable 미등록

증상

KeyError: Variable weather_decoding_api_key does not exist

원인

  • Airflow Variable이 등록되지 않은 상태에서 Variable.get() 호출

해결

  • Airflow UI 또는 CLI로 Variable 등록
airflow variablesset weather_decoding_api_key"<API_KEY>"

[7-2] 문제 2: 동적 start_date 설정으로 인한 설계 문제

start_date=pendulum.now(local_tz).subtract(hours=12)

문제점

  • DAG 파싱 시점마다 start_date가 변경됨
  • 스케줄 계산 및 catchup 동작이 예측 불가능해짐

교훈

  • start_date는 반드시 고정된 시점으로 설정
  • “최근 n시간만 처리” 같은 정책은 DAG 내부 로직으로 해결해야 함

[8] 현재 DAG의 핵심 특징 요약

  • Airflow Variable 기반 API Key 관리
  • data_interval_end 기준 시간 계산
  • API 미제공 구간은 실패가 아닌 정상 스킵 처리
  • JSON 파일 기반 결과 저장 구조
  • Spark/DB 파이프라인으로 확장 가능한 설계

업로드중..


[9] 정리 및 배운 점

  • Airflow는 단순 스케줄러가 아닌 데이터 파이프라인 오케스트레이터
  • 시간 개념(data_interval)을 정확히 이해하는 것이 핵심
  • 외부 API 연동 시 “실패”와 “정상적인 미제공”을 구분해야 함
  • Variable/Connection 관리가 운영 안정성에 직결됨
  • 처음부터 확장 가능한 구조로 DAG를 설계하면 이후 작업 비용이 크게 줄어듦

0개의 댓글