Operators
Operators
Operator가 무엇이고 어떻게 사용하는지 설명하는 문서예요. Operator는 미리 정의된 Task의 템플릿으로 Dag 안에서 선언적으로 정의할 수 있어요. Jinja 템플릿, params 예약 키워드, pre/post-execute 메서드, setup/teardown과의 차이까지 살펴볼게요.
출처: 문서
본문
Operator는 개념적으로 미리 정의된 Task의 템플릿이에요. 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에는 매우 광범위한 operator 세트가 있고, 일부는 core에 내장되거나 사전 설치된 providers에 있어요. core의 인기 있는 operator로는 다음이 있어요:
- BashOperator - bash 명령을 실행
- PythonOperator - 임의의 Python 함수를 호출
@task 데코레이터를 사용해 임의의 Python 함수를 실행할 수 있어요. 인자로 전달된 jinja 템플릿의 렌더링은 지원하지 않아요.
참고 (Note)
@task데코레이터는 인자에 템플릿 렌더링 없이 Python callable을 실행하는 PythonOperator의 Taskflow 동등물이에요.
모든 core operator 목록은 Core Operators and Hooks Reference를 참고하세요.
필요한 operator가 기본 Airflow 설치에 없으면, 방대한 커뮤니티 providers 세트의 일부로 찾을 수 있을 거예요. 여기서 인기 있는 operator로는 다음이 있어요:
- EmailOperator
- HttpOperator
- SQLExecuteQueryOperator
- DockerOperator
- HiveOperator
- S3FileTransformOperator
- PrestoToMySqlOperator
- SlackAPIOperator
하지만 훨씬 많아요 — 모든 커뮤니티 관리 operators, hooks, sensors와 transfers의 전체 목록은 우리의 providers packages 문서에서 볼 수 있어요.
참고 (Note)
Airflow 코드 내부에서는 Tasks와 Operators의 개념을 혼합해서 쓰는 경우가 많고, 대부분 상호 교환이 가능해요. 하지만 Task라고 말할 때는 Dag의 일반적인 "실행 단위"를 의미하고, Operator라고 말할 때는 로직이 모두 완성되어 인자만 필요로 하는 재사용 가능한 사전 제작 Task 템플릿을 의미해요.
Jinja 템플릿 (Jinja Templating)
Airflow는 Jinja Templating의 힘을 활용하며, 이는 macros와 함께 사용할 때 강력한 도구가 될 수 있어요.
예를 들어 데이터 구간의 시작을 환경 변수로 BashOperator를 사용하는 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 템플릿에서 사용할 수 있는 모든 변수를 포함하며, 템플릿 렌더링 관점에서 읽기 전용이에요 — 값을 접근하고 사용할 수 있지만, 수정이 태스크 실행 환경에는 영향을 미치지 않아요.
사용 가능한 context 변수의 전체 목록은 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"로 표시된 모든 파라미터에 템플릿을 사용할 수 있어요. 템플릿 치환은 operator의 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 문서에서 찾을 수 있어요.
일부 operator는 특정 접미사(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",
)
Airflow가 문자열을 명령이 아니라 파일 경로로 취급하므로 이것은 TemplateNotFound: cat script.sh로 실패해요. 값을 literal()로 감싸 Airflow가 이 값을 파일 참조로 취급하지 못하게 할 수 있어요. 이 접근 방식은 매크로와 파일 모두의 렌더링을 비활성화하고, 나머지 콘텐츠는 기본 템플릿 규칙을 유지하면서 선택된 중첩 필드에 적용할 수 있어요.
from airflow.sdk import literal
fixed_print_script = BashOperator(
task_id="fixed_print_script",
bash_command=literal("cat script.sh"),
)
2.8 버전에 추가됨:
literal()이 추가됨.
또는 Airflow가 값을 파일 참조로 취급하지 못하게 하려면 template_ext를 오버라이드할 수 있어요:
fixed_print_script = BashOperator(
task_id="fixed_print_script",
bash_command="cat script.sh",
)
fixed_print_script.template_ext = ()
필드를 네이티브 Python 객체로 렌더링하기 (Rendering Fields as Native Python Objects)
기본적으로 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
실제 dict를 얻고 싶다면 두 가지 해결책이 있어요. 첫 번째는 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 리터럴 규칙에 따라 값을 렌더링해요. 템플릿이 목록, dict, 숫자, 불리언을 생성해야 할 때 유용하지만,"42"처럼 숫자처럼 보이는 문자열도 정수42로 렌더링될 수 있다는 뜻이에요. 값이 문자열로 유지되어야 하는 태스크라면 기본 문자열 렌더링을 유지하고, callable 템플릿 필드를 사용하거나 명시적 인용을 추가하세요.
예약된 params 키워드 (Reserved params keyword)
Apache Airflow 2.2.0에서 params 변수는 Dag 직렬화 중에 사용돼요. 제3자 operator에서 그 이름을 사용하지 마세요. 환경을 업그레이드하고 다음 오류가 발생하면:
AttributeError: 'str' object has no attribute '__module__'
operator에서 params라는 이름을 변경하세요.
f-문자열과의 템플릿 충돌 (Templating Conflicts with f-strings)
템플릿 필드(예: 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- and post-execute methods)
pre_execute와 post_execute 메서드는 각각 operator가 실행되기 전과 후에 호출돼요.
예를 들어 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는 operator가 만든 임시 파일이나 디렉토리를 정리하는 데 사용할 수 있어요.
pre_execute와 post_execute 메서드에는 태스크 인스턴스의 context가 파라미터로 포함돼요.
pre-/post-execute와 setup/teardown의 차이 (Difference between pre-/post-execute and setup/teardown)
pre_execute와 post_execute 메서드는 개별 task instance 수준에서 operator가 실행되기 전과 후에 호출돼요. Setup과 teardown은 Dag run 내에서 여러 task instance가 실행되기 전/후에 setup 또는 cleanup 작업을 수행하는 데 사용되는 특수 태스크예요.