Executor (Executors)

Executor (Executors)

Executor는 task instance가 실제로 실행되는 메커니즘입니다. 모든 Executor는 공통된 API를 가지며, "플러그처럼" 교체할 수 있습니다. 다시 말해 설치 환경에 맞춰 Executor를 바꿔 끼울 수 있다는 뜻이에요.

Executor는 설정 파일[core] 섹션에 있는 executor 옵션으로 지정합니다.

내장 Executor는 이름으로 지정해요. 예를 들어:

 [core] executor = KubernetesExecutor

커스텀 또는 서드파티 Executor는 Executor 파이썬 클래스의 모듈 경로를 제공해 설정할 수 있습니다.

 [core] executor = my.custom.executor.module.ExecutorClass

참고: Airflow 설정에 대한 자세한 내용은 Setting Configuration Options을 참고하세요.

현재 설정된 Executor가 무엇인지 확인하고 싶다면 airflow config get-value core executor 명령을 사용하면 됩니다.

 $ airflow config get-value core executor LocalExecutor

Executor 유형

저장소의 트리 구조에서 태스크를 로컬에서(즉 scheduler 프로세스 안에서) 실행하는 Executor 유형은 하나뿐이지만, 커스텀 Executor를 작성해 비슷한 결과를 얻을 수도 있고, 태스크를 원격으로 실행(주로 worker 풀을 통해)하는 Executor도 있습니다. Airflow는 기본적으로 LocalExecutor로 설정되어 있는데, 이것은 로컬 Executor이며 실행 옵션 중 가장 단순한 선택이에요. 다만 LocalExecutor는 scheduler 프로세스 안에서 프로세스를 실행하므로 scheduler 성능에 영향을 줄 수 있습니다. LocalExecutor는 규모가 작은 단일 머신 프로덕션 설치에 쓰거나, 멀티 머신/클라우드 설치에는 원격 Executor 중 하나를 사용하면 됩니다.

로컬 Executor (Local Executors)

Airflow 태스크가 scheduler 프로세스 안에서 로컬로 실행됩니다.

장점: 사용이 매우 쉽고, 빠르며, 지연 시간이 매우 낮고, 설정 요구 사항도 적습니다.

단점: 기능이 제한적이고, Airflow scheduler와 리소스를 공유합니다.

예시:

원격 Executor (Remote Executors)

원격 Executor는 다시 두 가지 범주로 나뉩니다.

큐/배치 Executor

Airflow 태스크를 중앙 큐로 보내면 원격 worker가 태스크를 꺼내 실행합니다. worker는 보통 상시 실행(persistent)되며 여러 태스크를 동시에 처리하곤 해요.

장점: worker를 scheduler 프로세스에서 분리하므로 더 견고합니다. worker는 대형 호스트가 될 수 있어 많은 태스크(흔히 병렬로)를 처리할 수 있고 이는 비용 효율적입니다. 지연 시간도 비교적 낮을 수 있는데, worker가 항상 실행되도록 프로비저닝하면 큐에서 태스크를 즉시 가져올 수 있기 때문이에요.

단점: 공유 worker는 noisy neighbor(소음 이웃) 문제가 있어 태스크들이 공유 호스트의 리소스나 환경/시스템이 구성되는 방식을 두고 경쟁하게 됩니다. 또한 워크로드가 일정하지 않으면 비용이 커질 수 있어요. worker가 유휴 상태이거나 리소스가 과도하게 스케일링되거나, 스케일 업/다운을 관리해야 하기 때문입니다.

예시:

컨테이너화 Executor

Airflow 태스크가 컨테이너/파드 안에서 ad hoc 방식으로 실행됩니다. 각 태스크는 자체 컨테이너 환경에 격리되며, 이 환경은 Airflow 태스크가 큐에 들어갈 때 배포됩니다.

장점: 각 Airflow 태스크가 하나의 컨테이너에 격리되므로 noisy neighbor 문제가 없습니다. 실행 환경을 특정 태스크에 맞게 커스터마이즈할 수 있어요(시스템 라이브러리, 바이너리, 의존성, 리소스 양 등). worker가 태스크가 살아 있는 동안만 존재하므로 비용 효율적입니다.

단점: 컨테이너나 파드가 태스크를 시작하기 전에 배포되어야 하므로 시작 시 지연 시간이 있습니다. 많은 짧은/작은 태스크를 실행하면 비용이 커질 수 있어요. 관리할 worker는 없지만, 대신 Kubernetes 클러스터 같은 것을 관리해야 합니다.

예시:

참고: Airflow를 처음 쓰는 사용자라면 Local 또는 Remote Executor를 사용해 별도의 executor 프로세스를 실행해야 한다고 생각할 수 있어요. 이는 옳지 않습니다. Executor 로직은 scheduler 프로세스 안에서 실행되며, 선택된 Executor에 따라 태스크를 로컬로 실행할지 말지를 결정합니다.

여러 Executor 동시 사용

Airflow 2.10.0 버전부터 멀티 Executor 구성을 운영할 수 있습니다. 각 Executor에는 각자의 장단점이 있는데, 대개 지연 시간, 격리, 컴퓨팅 효율성 사이의 트레이드오프입니다(Executor 비교는 여기를 참고). 여러 Executor를 실행하면 사용 가능한 모든 Executor의 강점을 더 잘 활용하고 약점을 피할 수 있어요. 즉, 특정 태스크 집합에 그 가치와 이점이 가장 잘 맞는 특정 Executor를 사용할 수 있다는 뜻입니다.

구성

여러 Executor 구성은 단일 Executor 사용 사례와 동일한 구성 옵션을 사용하며, 쉼표로 구분된 목록 표기법을 활용해 여러 Executor를 지정합니다.

참고: 목록의 첫 번째 Executor(단독이든 다른 Executor와 함께든)는 2.10.0 이전 릴리스와 동일하게 동작합니다. 다시 말해 이 Executor가 환경의 기본 Executor가 돼요. 특정 Executor를 지정하지 않은 모든 Airflow Task 또는 Dag는 이 환경 수준의 Executor를 사용합니다. 목록의 다른 Executor는 모두 초기화되어 Airflow Task 또는 Dag에 지정되면 태스크를 실행할 준비가 됩니다. 이 구성 목록에 지정되지 않은 Executor는 태스크를 실행하는 데 사용할 수 없어요.

유효한 멀티 Executor 구성의 몇 가지 예시입니다.

 [core] executor = LocalExecutor
 [core] executor = LocalExecutor,CeleryExecutor
 [core] executor = KubernetesExecutor,my.custom.module.ExecutorClass

별칭 (Aliases)

태스크와 Dag에 Executor를 더 쉽게 지정하기 위해, Executor 구성은 별칭(alias)을 지원합니다. 그러면 Dag에서 이 별칭을 사용해 Executor를 가리킬 수 있어요(아래 참고).

별칭은 커스텀 Executor 모듈 경로와 내장 core Executor 모두에서 동작합니다.

 [core] executor = LocalExecutor,short_name:my.custom.module.ExecutorClass
 [core] executor = my_local_exec:LocalExecutor,my_celery_exec:CeleryExecutor

core Executor에 별칭을 붙이는 것은 Multi-Team Airflow를 운영할 때 특히 유용합니다. 태스크가 별칭으로 특정 인스턴스를 명시적으로 가리킬 수 있기 때문이에요.

 [core] executor = global_celery_exec:CeleryExecutor ;team1=team_celery_exec:CeleryExecutor

참고: 같은 Executor 클래스의 두 인스턴스를 사용하는 것은 멀티 팀 Airflow에서만 지원됩니다. 예를 들어 서로 다른 두 팀은 모두 CeleryExecutor를 사용할 수 있지만, 하나의 팀이 CeleryExecutor의 두 인스턴스를 사용할 수는 없어요. Executor는 전역으로 그리고 팀 안에서 동시에 사용할 수도 있습니다. 자세한 내용은 Multi-Team Airflow 문서를 참고하세요.

Dags와 Task 작성

참고: Dag가 구성되지 않은 Executor를 사용하도록 태스크를 지정하면 Dag는 파싱에 실패하고 Airflow UI에 경고 대화상자가 표시됩니다. 사용하려는 모든 Executor가 Airflow 구성 요소(scheduler, worker 등)를 실행하는 모든 호스트/컨테이너의 Airflow 구성에 지정되어 있는지 확인하세요.

태스크에 Executor를 지정하려면 Operator의 executor 파라미터를 사용합니다.

 BashOperator ( task_id = "hello_world" , executor = "LocalExecutor" , bash_command = "echo 'hello world!'" , )
 @task ( executor = "LocalExecutor" ) def hello_world (): print ( "hello world!" )

전체 Dag에 Executor를 지정하려면 기존 Airflow 메커니즘인 기본 인자(default arguments)를 사용합니다. 그러면 Dag의 모든 태스크가 지정된 Executor를 사용하게 돼요(특정 태스크가 명시적으로 재정의하지 않는 한).

 def hello_world (): print ( "hello world!" ) def hello_world_again (): print ( "hello world again!" ) with DAG ( dag_id = "hello_worlds" , default_args = { "executor" : "LocalExecutor" }, # Applies to all tasks in the Dag ) as dag : # All tasks will use the executor from default args automatically hw = hello_world () hw_again = hello_world_again ()

참고: 태스크는 실행되도록 구성된 Executor를 Airflow 데이터베이스에 저장합니다. 변경 사항은 Dag가 파싱될 때마다 반영됩니다.

모니터링

단일 Executor를 사용할 때 Airflow 메트릭은 2.9 이전 버전처럼 동작합니다. 하지만 여러 Executor가 구성되면 Executor 메트릭(executor.open_slots, executor.queued_slots, executor.running_tasks)이 구성된 각 Executor마다 발행되며, 메트릭 이름에 Executor 이름이 붙습니다(예: executor.open_slots.<executorclassname>).

로깅은 단일 Executor 사용 사례와 동일하게 동작합니다.

정적으로 코딩된 하이브리드 Executor (Statically-coded Hybrid Executors)

"정적으로 코딩된(statically coded)" Executor 두 개가 있었지만, Airflow 3.0.0부터 더 이상 지원되지 않습니다.

이 Executor들은 두 개의 서로 다른 Executor의 하이브리드입니다: LocalKubernetesExecutorCeleryKubernetesExecutor. 이들의 구현은 core Airflow에 고유하거나 본질적인 것이 아니에요. 이 하이브리드 Executor들은 대신 Task Instance의 queue 필드를 사용해 어떤 하위 Executor에서 실행할지 표시하고 유지합니다. 이는 queue 필드를 잘못 사용하는 것이며, 이 하이브리드 Executor를 사용할 때는 queue 필드를 원래 목적대로 사용할 수 없게 됩니다.

이런 Executor들은 또한 가능한 Executor 조합의 각 순열에 대해 새로운 "구체적인(concrete)" 클래스를 직접 만들어야 합니다. Executor가 늘어날수록 이는 감당하기 어려워지고 유지보수 오버헤드가 커져요. 실행하려는 Executor 조합에 대한 맞춤 코딩이 필요해서는 안 됩니다.

따라서 이런 유형의 Executor는 Airflow 3.0.0부터 더 이상 지원되지 않습니다. 대신 여러 Executor 동시 사용(Using Multiple Executors Concurrently) 기능을 사용하는 것이 좋습니다.

나만의 Executor 작성하기

모든 Airflow Executor는 공통 인터페이스를 구현하므로 플러그처럼 교체할 수 있고, 모든 Executor가 Airflow 안의 모든 기능과 통합에 접근할 수 있어요. 주로 Airflow scheduler가 이 인터페이스를 사용해 Executor와 상호작용하지만, 로깅이나 CLI 같은 다른 구성 요소도 마찬가지로 사용합니다. 공개 인터페이스는 BaseExecutor입니다. 가장 상세하고 최신의 인터페이스는 코드에서 살펴볼 수 있지만, 몇 가지 중요한 핵심을 아래에 정리했어요.

참고: Airflow 공개 인터페이스에 대한 자세한 내용은 Public Interface for Airflow 3.0+를 참고하세요.

커스텀 Executor를 작성하고 싶은 이유는 다음과 같습니다.

  • 컴퓨팅을 위한 특정 도구나 서비스 등, 특정 사용 사례에 맞는 Executor가 존재하지 않는 경우.
  • 선호하는 클라우드 공급자의 컴퓨팅 서비스를 활용하는 Executor를 사용하고 싶은 경우.
  • 나 또는 조직에만 제공되는, 태스크 실행을 위한 사설 도구/서비스를 가진 경우.

Workloads

Executor 맥락에서 workload는 Executor가 실행하는 기본 실행 단위(fundamental unit of execution)입니다. worker에서 Executor가 실행하는 개별 작업 또는 잡(job)을 나타냅니다. 예를 들어 Airflow 태스크에 캡슐화된 사용자 코드를 worker에서 실행할 수 있어요.

예시:

 ExecuteTask ( token = "mock" , ti = TaskInstanceDTO ( id = UUID ( "4d828a62-a417-4936-a7a6-2b3fabacecab" ), task_id = "mock" , dag_id = "mock" , run_id = "mock" , try_number = 1 , dag_version_id = UUID ( "4d828a62-a417-4936-a7a6-2b3fabacecab" ), map_index =- 1 , pool_slots = 1 , queue = "default" , priority_weight = 1 , executor_config = None , ), dag_rel_path = PurePosixPath ( "mock.py" ), bundle_info = BundleInfo ( name = "n/a" , version = "no matter" ), log_path = "mock.log" , type = "ExecuteTask" , )

중요한 BaseExecutor 메서드

이 메서드들은 나만의 Executor를 구현하기 위해 반드시 재정의할 필요는 없지만, 알아 두면 유용합니다.

  • heartbeat: Airflow scheduler Job 루프가 주기적으로 Executor의 heartbeat를 호출합니다. 이것이 Airflow scheduler와 Executor 사이의 주요 상호작용 지점 중 하나에요. 이 메서드는 일부 메트릭을 갱신하고, 새로 큐에 들어온 태스크가 실행되도록 트리거하며, 실행 중/완료된 태스크의 상태를 갱신합니다.
  • queue_workload: Airflow Executor가 이 BaseExecutor 메서드를 호출해 Executor가 실행할 태스크를 제공합니다. BaseExecutor는 단순히 workload들(위 섹션을 참고)을 Executor 안에서 실행할 큐에 들어온 workload의 내부 목록에 추가합니다. 저장소에 존재하는 모든 Executor가 이 메서드를 사용합니다.
  • get_event_buffer: Airflow scheduler가 이 메서드를 호출해 Executor가 실행 중인 TaskInstance의 현재 상태를 가져옵니다.
  • has_task: scheduler가 이 BaseExecutor 메서드를 사용해 Executor가 특정 태스크 인스턴스를 이미 큐에 넣었거나 실행 중인지 판단합니다.
  • send_callback: Executor에 구성된 sink로 콜백을 보냅니다.

반드시 구현해야 하는 메서드

나만의 Executor가 Airflow에서 지원되도록 하려면 최소한 다음 메서드를 재정의해야 합니다.

  • sync: sync는 Executor heartbeat 중에 주기적으로 호출됩니다. 이 메서드를 구현해 Executor가 알고 있는 태스크의 상태를 갱신합니다. 선택적으로 scheduler로부터 받은 큐에 들어온 태스크를 실행하려고 시도할 수 있어요.
  • execute_async: workload를 비동기로 실행합니다. 이 메서드는 scheduler가 주기적으로 실행하는 Executor heartbeat 중에 (몇 겹의 계층을 거쳐) 호출됩니다. 실제로는 이 메서드가 내부 또는 외부의 태스크 실행 큐에 태스크를 넣는 경우가 많아요(예: KubernetesExecutor). 하지만 태스크를 직접 실행할 수도 있습니다(예: LocalExecutor). 이는 Executor에 따라 달라집니다.
  • _process_workloads: queue_workload를 통해 큐에 들어온 workload 목록을 처리합니다. 이 메서드는 Executor heartbeat 중에 호출되며 Executor가 workload 실행을 어떻게 처리할지 정의합니다(예: worker에 큐잉, 외부 시스템에 제출 등).

구현할 수 있는 선택적 인터페이스 메서드

다음 메서드들은 동작하는 Airflow Executor를 만들기 위해 재정의할 필요는 없습니다. 하지만 구현하면 강력한 기능과 안정성을 얻을 수 있어요.

  • start: Airflow scheduler 잡이 Executor 객체를 초기화한 후 이 메서드를 호출합니다. Executor가 필요한 추가 설정을 여기서 완료할 수 있어요.
  • end: Airflow scheduler 잡이 teardown하면서 이 메서드를 호출합니다. 실행 중인 잡을 끝내는 데 필요한 동기적 정리를 여기서 해야 합니다.
  • terminate: Executor를 더 강력하게 중지합니다. 완료를 동기적으로 기다리는 대신 실행 중(in-flight)인 태스크를 죽이거나 중지하기도 합니다.
  • try_adopt_task_instances: 버려진 태스크(예: 죽은 scheduler 잡에서)는 이 메서드를 통해 Executor에 제공되어 인계(adopt)되거나 다르게 처리됩니다. 인계할 수 없는 태스크(기본적으로 BaseExecutor는 모든 태스크를 인계할 수 없다고 가정)는 반환되어야 합니다.
  • get_cli_commands: 이 메서드를 구현하면 Executor가 사용자에게 CLI 명령을 제공(vend)할 수 있어요. 자세한 내용은 아래의 CLI 섹션을 참고하세요.
  • get_task_log: 이 메서드를 구현하면 Executor가 Airflow 태스크 로그에 로그 메시지를 제공할 수 있습니다. 자세한 내용은 아래의 Logging 섹션을 참고하세요.

호환성 속성 (Compatibility Attributes)

BaseExecutor 클래스 인터페이스에는 Airflow core 코드가 Executor가 호환되는 기능을 확인하는 데 사용하는 속성 집합이 있습니다. 나만의 Airflow Executor를 작성할 때는 사용 사례에 맞게 이 속성들을 정확히 설정하세요. 각 속성은 단순히 기능을 활성화/비활성화하거나 Executor가 지원/미지원함을 나타내는 불리언입니다.

  • supports_pickling: Executor가 실행 전에 (파일 시스템에서 Dag 정의를 읽는 대신) 데이터베이스에서 pickled Dag를 읽는 것을 지원하는지 여부.
  • sentry_integration: Executor가 Sentry를 지원한다면, 통합을 만드는 callable의 임포트 경로여야 합니다. 예를 들어 CeleryExecutor는 이 값을 "sentry_sdk.integrations.celery.CeleryIntegration"으로 설정합니다.
  • is_local: Executor가 원격인지 로컬인지 여부. 위의 Executor 유형 섹션을 참고하세요.
  • is_single_threaded: Executor가 단일 스레드인지 여부. 어떤 데이터베이스 백엔드를 지원하는지와 특히 관련이 있습니다. 단일 스레드 Executor는 SQLite를 포함한 어떤 백엔드에서도 실행할 수 있어요.
  • is_production: Executor를 프로덕션 목적으로 사용해야 하는지 여부. 프로덕션 준비가 되지 않은 Executor를 사용하면 사용자에게 UI 메시지가 표시됩니다.
  • serve_logs: Executor가 로그 제공(serving logs)을 지원하는지 여부. Logging for Tasks를 참고하세요.

CLI

중요: Airflow 3.2.0부터 auth manager와 executor 같은 core 확장을 관리하는 provider 수준 CLI 명령을 사용할 수 있습니다. provider 수준 CLI 명령을 구현하면 필요하지 않을 때 무거운 임포트를 피해 CLI 시작 시간을 줄일 수 있어요. 구현 안내는 provider-level CLI를 참고하세요.

Executor는 get_cli_commands 메서드를 구현해 airflow 커맨드라인 도구에 포함될 CLI 명령을 제공할 수 있습니다. 예를 들어 CeleryExecutorKubernetesExecutor 같은 Executor가 이 메커니즘을 사용해요. 이 명령들은 필요한 worker 설정, 환경 초기화, 또는 다른 구성 설정에 사용할 수 있습니다. 명령은 현재 구성된 Executor에 대해서만 제공됩니다. Executor에서 CLI 명령 제공을 구현하는 의사 코드 예시는 아래와 같습니다.

 @staticmethod def get_cli_commands () -> list [ GroupCommand ]: sub_commands = [ ActionCommand ( name = "command_name" , help = "Description of what this specific command does" , func = lazy_load_command ( "path.to.python.function.for.command" ), args = (), ), ] return [ GroupCommand ( name = "my_cool_executor" , help = "Description of what this group of commands do" , subcommands = sub_commands , ), ]

참고: 현재 Airflow 명령 네임스페이스에는 엄격한 규칙이 없습니다. CLI 명령 이름이 다른 Airflow Executor나 구성 요소와 충돌하지 않도록 충분히 고유하게 짓는 것은 개발자 몫입니다.

참고: 새 Executor를 만들거나 기존 Executor를 갱신할 때는 모듈 레벨에서 비싼 연산/코드를 임포트하거나 실행하지 마세요. Executor 클래스는 여러 곳에서 임포트되며, 임포트가 느리면 특히 CLI 명령에서 Airflow 환경의 성능에 부정적인 영향을 줍니다.

로깅 (Logging)

Executor는 get_task_logs 메서드를 구현해 Airflow 태스크 로그에 포함될 로그 메시지를 제공할 수 있습니다. 실행 환경이 태스크 실패 시 추가 맥락을 가진 경우(이 실패가 Airflow 태스크 코드 때문이 아니라 실행 환경 자체 때문일 수 있음) 유용해요. 실행 환경의 setup/teardown 로깅을 포함하는 데에도 도움이 됩니다. KubernetesExecutor는 이 기능을 활용해 특정 Airflow 태스크를 실행한 pod의 로그를 가져와 그 Airflow 태스크의 로그에 표시합니다. Executor에서 태스크 로그 제공을 구현하는 의사 코드 예시는 아래와 같습니다.

 def get_task_log ( self , ti : TaskInstance , try_number : int ) -> tuple [ list [ str ], list [ str ]]: messages = [] log = [] try : res = helper_function_to_fetch_logs_from_execution_env ( ti , try_number ) for line in res : log . append ( remove_escape_codes ( line . decode ())) if log : messages . append ( "Found logs from execution environment!" ) except Exception as e : # No exception should cause task logs to fail messages . append ( f "Failed to find logs from execution environment: { e } " ) return messages , [ " \\n " . join ( log )]