getUltraSrtNcst)를 주기적으로 호출data_interval_end를 기준으로 조회 시간 계산data_interval_end를 기준 시점으로 사용now = data_interval_end.in_timezone("Asia/Seoul")
이 방식의 장점:
api_key = Variable.get("weather_decoding_api_key")
이유:
include 디렉토리 사용/usr/local/airflow/include/ultra_srt_weather/
└── ncst_YYYYMMDD_HH00.json
이유:
fetch_task = PythonOperator(
task_id="fetch_ncst_step",
python_callable=fetch_ultra_srt_ncst
)
schedule="15 * * * *"
catchup=True
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
)
KeyError: Variable weather_decoding_api_key does not exist
Variable.get() 호출airflow variablesset weather_decoding_api_key"<API_KEY>"
start_date=pendulum.now(local_tz).subtract(hours=12)
start_date는 반드시 고정된 시점으로 설정data_interval_end 기준 시간 계산
data_interval)을 정확히 이해하는 것이 핵심