카테고리 없음
[Airflow] Custom 오퍼레이터 (2)
by imkus
2025. 1. 13.
1. Custom 오퍼레이터를 사용하는 이유
- 만약 Custom 오퍼레이터를 만들지 않았다면?
- 개발자마다 각자 서울 공공데이터 데이터셋 추출/저장하는 파이썬 파일을 만들어 Python Operator를 이용해 개발했을 것
- 비슷한 동작을 하는 파이썬 파일이 관리되지 않은 채 수십 개 만들어지면 그 자체로 비효율 발생
- 특정 기능을 하는 모듈을 만들어 놓고 상세 조건은 파라미터로 받도록하여 모듈을 재사용할 수 있도록 유도
2. Custom 오퍼레이터 사용 예제코드
'''
plugins/operators/seoul_api_to_csv_operator.py
'''
from airflow.models.baseoperator import BaseOperator
from airflow.hooks.base import BaseHook
import pandas as pd
class SeoulApiToCsvOperator(BaseOperator):
template_fields = ('endpoint', 'path', 'file_name', 'base_dt')
def __init__(self, dataset_nm, path, file_name, base_dt=None, **kwargs):
super().__init__(**kwargs)
self.http_conn_id = 'openapi.seoul.go.kr'
self.path = path
self.file_name = file_name
self.endpoint = '{{var.value.apikey_openapi_seoul_go_kr}}/json/' + dataset_nm
self.base_dt = base_dt
def execute(self, context):
import os
connection = BaseHook.get_connection(self.http_conn_id)
self.base_url = f'http://{connection.host}:{connection.port}/{self.endpoint}'
total_row_df = pd.DataFrame()
start_row = 1
end_row = 1000
while True:
self.log.info(f'시작:{start_row}')
self.log.info(f'끝:{end_row}')
row_df = self._call_api(self.base_url, start_row, end_row)
total_row_df = pd.concat([total_row_df, row_df])
if len(row_df) < 1000:
break
else:
start_row = end_row + 1
end_row += 1000
if not os.path.exists(self.path):
os.system(f'mkdir -p {self.path}')
total_row_df.to_csv(self.path + '/' + self.file_name, encoding='utf-8', index=False)
def _call_api(self, base_url, start_row, end_row):
import requests
import json
headers = {'Content-Type':'application/json',
'charset': 'utf-8',
'Accept': '*/*'
}
request_url = f'{base_url}/{start_row}/{end_row}/'
if self.base_dt is not None:
request_url = f'{base_url}/{start_row}/{end_row}/{self.base_dt}'
response = requests.get(request_url, headers)
contents = json.loads(response.text)
key_nm = list(contents.keys())[0]
row_data = contents.get(key_nm).get('row')
row_df = pd.DataFrame(row_data)
return row_df
'''
dags_seoul_api_corona.py
'''
from operators.seoul_api_to_csv_operator import SeoulApiToCsvOperator
from airflow import DAG
import pendulum
with DAG(
dag_id = 'dags_seoul_api_corona',
schedule = '0 7 * * *',
start_date = pendulum.datetime(2023, 4, 1, tz='Asia/Seoul'),
catchup = False
) as dag:
'''서울시 코로나19 확진자 발생동향'''
tb_corona19_count_status = SeoulApiToCsvOperator(
task_id = 'tb_corona19_count_status',
dataset_nm = 'TBCorona19CountStatus',
path='/opt/airflow/files/TbCorona19CountStatus/{{data_interval_end.in_timezone("Asia/Seoul") | ds_nodash }}',
file_name = 'TbCorona19CountStatus.csv'
)
'''서울시 코로나19 백신 예방접종 현황'''
tv_corona19_vaccine_stat_new = SeoulApiToCsvOperator(
task_id = 'tv_corona19_vaccine_stat_new',
dataset_nm = 'tvCorona19VaccinestatNew',
path='/opt/airflow/files/tvCorona19VaccinestatNew/{{data_interval_end.in_timezone("Asia/Seoul") | ds_nodash }}',
file_name = 'tvCorona19VaccinestatNew.csv'
)
tb_corona19_count_status >> tv_corona19_vaccine_stat_new