본문 바로가기
카테고리 없음

[Airflow] Custom 오퍼레이터 (2)

by imkus 2025. 1. 13.

1. Custom 오퍼레이터를 사용하는 이유

  • 만약 Custom 오퍼레이터를 만들지 않았다면?
    • 개발자마다 각자 서울 공공데이터 데이터셋 추출/저장하는 파이썬 파일을 만들어 Python Operator를 이용해 개발했을 것
    • 비슷한 동작을 하는 파이썬 파일이 관리되지 않은 채 수십 개 만들어지면 그 자체로 비효율 발생
  • 특정 기능을 하는 모듈을 만들어 놓고 상세 조건은 파라미터로 받도록하여 모듈을 재사용할 수 있도록 유도
    • Custom 오퍼레이터 개발

 

2. Custom 오퍼레이터 사용 예제코드

  • 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

 

 

  • Custom 오퍼레이터를 사용한 DAG 생성
'''
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