오퍼레이터 (Operators)
오퍼레이터 (Operators)
오퍼레이터(Operator)는 개념적으로 미리 정의된 태스크를 위한 템플릿이에요. DAG 안에서 선언적으로 정의만 하면 되죠.
with DAG("my-dag") as dag:
ping = HttpOperator(endpoint="http://example.com/update/")
email = EmailOperator(to="[email protected]", subject="Update complete")
ping >> email
Airflow는 아주 방대한 오퍼레이터 세트를 제공해요. 코어에 기본 내장된 것도 있고, 사전 설치된 프로바이더(provider)에 포함된 것도 있죠. 코어에서 많이 쓰는 오퍼레이터를 몇 개 꼽아볼게요.
BashOperator— bash 명령을 실행해요PythonOperator— 임의의 Python 함수를 호출해요@task데코레이터 — 임의의 Python 함수를 실행해요. 단, 인자로 전달된 Jinja 템플릿 렌더링은 지원하지 않아요
Note
인자에 템플릿 렌더링이 필요 없는 Python callable을 실행할 때는 클래식한
PythonOperator보다@task데코레이터를 권장해요.
코어 오퍼레이터 전체 목록은 Core Operators and Hooks Reference에서 볼 수 있어요.
Airflow에 기본으로 설치되어 있지 않은 오퍼레이터가 필요하다면, 우리가 만든 방대한 커뮤니티 프로바이더 세트에서 찾을 수 있을 거예요. 여기서 인기 있는 오퍼레이터는 다음과 같아요.
EmailOperatorHttpOperatorSQLExecuteQueryOperatorDockerOperatorHiveOperatorS3FileTransformOperatorPrestoToMySqlOperatorSlackAPIOperator
그리고 이보다 훨씬, 훨씬 많아요. 커뮤니티가 관리하는 모든 오퍼레이터, 훅(hook), 센서(sensor), 트랜스퍼(transfer)의 전체 목록은 프로바이더 패키지 문서에서 확인할 수 있어요.
Note
Airflow 코드 안에서는 태스크(Task)와 오퍼레이터(Operator) 개념을 섞어 쓰는 경우가 많아서, 대체로 서로 바꿔 써도 돼요. 다만 얘기할 때 태스크는 DAG의 일반적인 '실행 단위'를, 오퍼레이터는 로직이 전부 구현되어 있어서 인자만 채워 넣으면 되는, 재사용 가능한 미리 만들어진 태스크 템플릿을 뜻해요.
Jinja 템플릿
Airflow는 Jinja 템플릿의 힘을 활용해요. 이를 매크로와 함께 쓰면 아주 강력한 도구가 되죠.
예를 들어 BashOperator를 써서 데이터 구간(data interval)의 시작을 환경 변수로 Bash 스크립트에 넘기고 싶다고 해볼게요.
# The start of the data interval as YYYY-MM-DD
date = "{{ ds }}"
t = BashOperator(
task_id="test_env",
bash_command="/tmp/test.sh ",
dag=dag,
env={"DATA_INTERVAL_START": date},
)
여기서 {{ds}}는 템플릿 변수예요. BashOperator의 env 파라미터가 Jinja로 템플릿 처리되기 때문에, 데이터 구간 시작 날짜가 Bash 스크립트 안에서 DATA_INTERVAL_START라는 이름의 환경 변수로 들어가게 돼요.
Jinja 템플릿보다 Python이 더 읽기 쉬운 상황이라면, 대신 callable을 넘길 수도 있어요. 그 callable은 context와 jinja_env라는 두 개의 명명된 인자를 받아야 해요.
context 파라미터는 현재 태스크 실행에 대한 런타임 정보를 제공하는 Airflow의 Context 객체예요. 그 내용물은 Python 표준 dict 문법으로 접근할 수 있어요. 여기에는 Jinja 템플릿에서 쓸 수 있는 모든 변수가 포함되어 있고, 템플릿 렌더링 관점에서 보면 읽기 전용이에요. 값에 접근해 쓰는 것은 가능하지만, 그렇게 수정한 내용이 태스크 실행 환경에는 영향을 주지 않아요.
쓸 수 있는 컨텍스트 변수의 전체 목록은 Templates reference를 참고하세요.
from typing import TYPE_CHECKING
if TYPE_CHECKING:
import jinja2
from airflow.sdk import Context
def build_complex_command(context: Context, jinja_env: jinja2.Environment) -> str:
# Access runtime information from the context dictionary
task_id = context["ti"].task_id
execution_date = context["ds"]
with open("file.csv") as f:
return do_complex_things(f, task_id, execution_date)
t = BashOperator(
task_id="complex_templated_echo",
bash_command=build_complex_command,
dag=dag,
)
각 템플릿 필드는 한 번만 렌더링되기 때문에, callable의 반환값은 다시 렌더링을 거치지 않아요. 그래서 callable은 템플릿을 직접 렌더링해야 해요. 이렇게 하려면 현재 태스크에서 render_template()을 호출하면 돼요.
def build_complex_command(context: Context, jinja_env: jinja2.Environment) -> str:
with open("file.csv") as f:
data = do_complex_things(f)
return context["task"].render_template(data, context, jinja_env)
문서에서 'templated'라고 표시된 파라미터라면 모두 템플릿 처리를 쓸 수 있어요. 템플릿 치환은 오퍼레이터의 pre_execute 함수가 호출되기 바로 직전에 일어나요.
중첩된 필드에도 템플릿 처리를 쓸 수 있어요. 단, 그 중첩 필드가 속한 구조에서 templated로 표시되어 있어야 해요. template_fields 프로퍼티에 등록된 필드는 템플릿 치환 대상이 되는데, 아래 예시의 path 필드처럼요.
class MyDataReader:
template_fields: Sequence[str] = ("path",)
def __init__(self, my_path):
self.path = my_path
# [additional code here...]
t = PythonOperator(
task_id="transform_data",
python_callable=transform_data,
op_args=[MyDataReader("/tmp/{{ ds }}/my_file")],
dag=dag,
)
Note
template_fields프로퍼티는 클래스 변수이며 항상Sequence[str]타입(즉 문자열 리스트 또는 튜플)임이 보장돼요.
아주 깊이 중첩된 필드도 치환할 수 있어요. 중간에 있는 모든 필드가 템플릿 필드로 표시되어 있다면 말이죠.
class MyDataTransformer:
template_fields: Sequence[str] = ("reader",)
def __init__(self, my_reader):
self.reader = my_reader
# [additional code here...]
class MyDataReader:
template_fields: Sequence[str] = ("path",)
def __init__(self, my_path):
self.path = my_path
# [additional code here...]
t = PythonOperator(
task_id="transform_data",
python_callable=transform_data,
op_args=[MyDataTransformer(MyDataReader("/tmp/{{ ds }}/my_file"))],
dag=dag,
)
DAG를 만들 때 Jinja Environment에 커스텀 옵션을 넘길 수도 있어요. 흔한 용례 중 하나가 Jinja가 템플릿 문자열의 끝에 있는 개행을 떨어뜨리지 않게 하는 거예요.
my_dag = DAG(
dag_id="my-dag",
jinja_environment_kwargs={
"keep_trailing_newline": True,
# some other jinja2 Environment options here
},
)
사용 가능한 모든 옵션은 Jinja 문서에서 확인할 수 있어요.
또 어떤 오퍼레이터들은 template_ext에 정의된 특정 접미사로 끝나는 문자열을 필드를 렌더링할 때 파일 참조로 간주하기도 해요. 이 기능은 스크립트나 쿼리를 DAG 코드에 직접 넣지 않고 파일에서 불러올 때 유용하죠.
예를 들어 여러 줄짜리 bash 스크립트를 실행하는 BashOperator를 생각해볼게요. 그러면 script.sh 파일을 읽어서 그 내용물을 bash_command 값으로 써요.
run_script = BashOperator(
task_id="run_script",
bash_command="script.sh",
)
기본적으로 이런 방식으로 제공되는 경로는 DAG 폴더를 기준으로 한 상대 경로여야 해요(Jinja 템플릿 검색 경로의 기본값이니까요). 하지만 DAG에 template_searchpath 인자를 설정하면 경로를 추가로 더할 수 있어요.
어떤 경우에는 문자열을 템플릿 처리에서 제외하고 그대로 쓰고 싶을 수도 있어요. 다음 태스크를 살펴볼게요.
print_script = BashOperator(
task_id="print_script",
bash_command="cat script.sh",
)
이건 TemplateNotFound:catscript.sh 오류로 실패해요. Airflow가 이 문자열을 명령어가 아니라 파일 경로로 취급하기 때문이죠. 이 값을 파일 참조로 취급하지 않게 하려면 literal()로 감싸면 돼요. 이 방식은 매크로와 파일 렌더링을 모두 비활성화하며, 선택한 중첩 필드에만 적용하고 나머지 내용에는 기본 템플릿 규칙을 유지할 수 있어요.
from airflow.sdk import literal
fixed_print_script = BashOperator(
task_id="fixed_print_script",
bash_command=literal("cat script.sh"),
)
버전 2.8에서 추가됨: literal()이 추가되었어요.
대신 값을 파일 참조로 취급하지 않게 하려면 template_ext를 오버라이드할 수도 있어요.
fixed_print_script = BashOperator(
task_id="fixed_print_script",
bash_command="cat script.sh",
)
fixed_print_script.template_ext = ()
필드를 네이티브 Python 객체로 렌더링하기
기본적으로 template_fields 안의 모든 Jinja 템플릿은 문자열로 렌더링돼요. 그런데 항상 그게 바람직한 건 아니에요. 예를 들어 extract 태스크가 {"1001":301.27,"1002":433.21,"1003":502.22} 같은 딕셔너리를 XCom으로 푸시한다고 해볼게요.
@task(task_id="extract")
def extract():
data_string = '{"1001": 301.27, "1002": 433.21, "1003": 502.22}'
return json.loads(data_string)
extract에 의존하는 태스크가 있다면, order_data 인자에는 "{'1001':301.27,'1002':433.21,'1003':502.22}"라는 문자열이 전달돼요.
def transform(order_data):
total_order_value = sum(order_data.values()) # Fails because order_data is a str :(
return {"total_order_value": total_order_value}
transform = PythonOperator(
task_id="transform",
op_kwargs={"order_data": "{{ ti.xcom_pull('extract') }}"},
python_callable=transform,
)
extract() >> transform
실제 딕셔너리를 얻고 싶다면 해결책이 두 가지 있어요. 첫 번째는 callable을 쓰는 거예요.
def render_transform_op_kwargs(context, jinja_env):
order_data = context["ti"].xcom_pull("extract")
return {"order_data": order_data}
transform = PythonOperator(
task_id="transform",
op_kwargs=render_transform_op_kwargs,
python_callable=transform,
)
두 번째로, Jinja에 네이티브 Python 객체를 렌더링하라고 지시할 수도 있어요. 이렇게 하려면 DAG에 render_template_as_native_obj=True를 전달하면 돼요. 그러면 Airflow가 기본 SandboxedEnvironment 대신 NativeEnvironment를 사용해요.
with DAG(
dag_id="example_template_as_python_object",
schedule=None,
start_date=pendulum.datetime(2021, 1, 1, tz="UTC"),
catchup=False,
render_template_as_native_obj=True,
):
transform = PythonOperator(
task_id="transform",
op_kwargs={"order_data": "{{ ti.xcom_pull('extract') }}"},
python_callable=transform,
)
Note
NativeEnvironment는 Python 리터럴 규칙에 따라 값을 렌더링해요. 템플릿이 리스트, 딕셔너리, 숫자, 불리언을 만들어야 할 때 유용하죠. 다만 그만큼"42"처럼 숫자처럼 보이는 문자열이 정수42로 렌더링될 수도 있어요. 값이 문자열로 남아야 하는 태스크라면 기본 문자열 렌더링을 유지하거나, callable 템플릿 필드를 쓰거나, 명시적인 따옴표를 추가하세요.
예약된 params 키워드
Apache Airflow 2.2.0에서 params 변수는 DAG 직렬화 중에 사용돼요. 서드파티 오퍼레이터에서는 그 이름을 쓰지 마세요. 환경을 업그레이드했는데 다음 오류가 난다면,
AttributeError: 'str' object has no attribute '__module__'
오퍼레이터의 이름을 params에서 바꾸세요.
f-스트링과의 템플릿 충돌
템플릿 필드(예: BashOperator의 bash_command)를 위한 문자열을 Python f-스트링으로 만들 때는, f-스트링 보간과 Jinja 템플릿 문법 사이의 상호작용에 주의해야 해요. 둘 다 중괄호({})를 쓰니까요.
Python f-스트링은 겹중괄호({{와 }})를 리터럴 단일 중괄호({와 })의 이스케이프 시퀀스로 해석해요. 반면 Jinja는 템플릿 변수를 나타낼 때 겹중괄호({{variable}})를 쓰죠.
f-스트링으로 정의한 문자열 안에 Jinja 템플릿 표현식(예: {{ds}})을 리터럴로 넣어서, 나중에 Airflow의 Jinja 엔진이 처리하게 하려면, f-스트링용 중괄호를 다시 한 번 더 이중화해서 이스케이프해야 해요. 즉 중괄호를 네 개 써야 한다는 뜻이에요.
t1 = BashOperator(
task_id="fstring_templating_correct",
bash_command=f"echo Data interval start: {{{{ ds }}}}",
dag=dag,
)
python_var = "echo Data interval start:"
t2 = BashOperator(
task_id="fstring_templating_simple",
bash_command=f"{python_var} {{{{ ds }}}}",
dag=dag,
)
이렇게 하면 f-스트링 처리를 거친 결과에 Jinja가 필요로 하는 리터럴 겹중괄호가 남게 되고, Airflow가 실행 전에 제대로 템플릿 처리할 수 있어요. 이걸 지키지 않으면 초보자에게 아주 흔한 문제인데, DAG 파싱 중 오류가 나거나, 기대한 대로 템플릿 처리가 이뤄지지 않아 런타임에 예상 밖의 동작이 발생할 수 있어요.
pre/post-execute 메서드
pre_execute와 post_execute 메서드는 각각 오퍼레이터가 실행되기 전과 후에 호출돼요.
예를 들어 pre_execute 메서드를 쓰면 태스크를 실행할지 말지를 우아하게 결정할 수 있어요.
def _check_skipped(context: Any) -> None:
"""
Check if a given task instance should be skipped if the `tasks_to_skip` Airflow Variable is a list that contains the task id.
"""
tasks_to_skip = Variable.get("tasks_to_skip", deserialize_json=True)
if context["task"].task_id in on_ice:
raise AirflowSkipException("Task instance configured to be skipped by `tasks_to_skip` variable.")
@task(pre_execute=_check_skipped)
def test():
"""
This task will be skipped if the `tasks_to_skip` Airflow Variable is set and contains the task id.
"""
...
post_execute는 오퍼레이터가 만든 임시 파일이나 디렉터리를 정리하는 데 쓸 수 있어요.
pre_execute와 post_execute 메서드는 인자로 태스크 인스턴스의 컨텍스트를 받아요.
pre/post-execute와 setup/teardown의 차이
pre_execute와 post_execute 메서드는 개별 태스크 인스턴스 수준에서 오퍼레이터 실행 전후에 호출돼요. 반면 setup과 teardown은 DAG 실행(run) 안에서 여러 태스크 인스턴스가 실행되기 전후에 설정 또는 정리 작업을 수행하는 특수 태스크예요.
Apache Airflow, Apache, Airflow 로고, Apache 로고는 The Apache Software Foundation의 등록 상표 또는 상표입니다. 그 외 모든 제품·브랜드 이름은 각 소유자의 상표입니다.